first commit
This commit is contained in:
85
code/utils/vdes_curd.py
Normal file
85
code/utils/vdes_curd.py
Normal file
@@ -0,0 +1,85 @@
|
||||
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])
|
||||
Reference in New Issue
Block a user