import logging from django.db import transaction, IntegrityError from time import perf_counter from system.models import ShipLatestStatus logger = logging.getLogger(__name__) def safe_bulk_upsert(ship_data_list): start_time = perf_counter() created = updated = 0 try: # 验证数据格式 valid_data = [] for data in ship_data_list: if 'mmsi' not in data: logger.warning("Missing MMSI in data: %s", data) continue valid_data.append(data) if not valid_data: return 0, 0 # 执行批量操作 with transaction.atomic(): current_mmsi = {item['mmsi'] for item in valid_data} existing_map = ShipLatestStatus.objects.in_bulk( current_mmsi, field_name='mmsi' ) to_create = [] to_update = [] for data in valid_data: mmsi = data['mmsi'] if mmsi in existing_map: ship = existing_map[mmsi] for field, value in data.items(): setattr(ship, field, value) to_update.append(ship) else: to_create.append(ShipLatestStatus(**data)) if to_create: created = len(ShipLatestStatus.objects.bulk_create(to_create)) if to_update: update_fields = [f.name for f in ShipLatestStatus._meta.fields if f.name not in ['id', 'mmsi']] ShipLatestStatus.objects.bulk_update(to_update, update_fields) updated = len(to_update) except IntegrityError as e: logger.error("Integrity error in bulk upsert: %s", str(e)) # 这里可以添加重试或更精细的错误处理 raise except Exception as e: logger.exception("Unexpected error in bulk upsert") raise elapsed = (perf_counter() - start_time) * 1000 logger.info(f"Processed {len(valid_data)} ships in {elapsed:.2f}ms: " f"{created} created, {updated} updated") return created, updated def get_nearby_ships(latitude: float, longitude: float, radius_nm: float = 20.0): # 将海里转换为公里(Haversine公式使用公里) radius_km = radius_nm * 1.852 # 使用原生SQL计算 return ShipLatestStatus.objects.raw(""" SELECT *, (6371 * acos( cos(radians(%s)) * cos(radians(latitude)) * cos(radians(longitude) - radians(%s)) + sin(radians(%s)) * sin(radians(latitude)) )) AS distance FROM ship_latest_status HAVING distance <= %s ORDER BY distance """, [latitude, longitude, latitude, radius_km])