fix(sync): 修复 Oracle id 重排导致的 sync_videos/sync_people 1062

Oracle 库重建/重排后 videos/people 的 id 跳变(2303 段→3408/4600+ 段),
旧 upsert 按 Oracle id 作主键插入时与已存在同 filename/label 行在 UNIQUE 键
上二次冲突报 1062,整批中止、游标不推进。

- upsert_sync_videos: 改为以 filename 业务键去重,命中就地 UPDATE 保留 NAS 原
  id(防 sync_events/model_calls/identity_map 的 video_id 外键失效),Oracle id
  落 oracle_id 列溯源;未命中插入优先用 Oracle id 对齐子表引用,主键冲突回退自增
- upsert_sync_people: 同模式,以 label 业务键去重(people.id 无外键引用)
- ddl.sql: sync_videos/sync_people 的 id 改为 NAS 本地自增主键 + 新增 oracle_id 列
- NAS 现网已迁移(加列/回填/改自增)并重启验证:单批补拉 videos+348 events+1017
  people+9 model_calls+500 identity_map+238,游标 09-02 13:55:21 → 09-03 08:51:11
This commit is contained in:
ericwyuan
2026-09-03 09:00:50 +08:00
parent 2ad0472578
commit 0238030128
2 changed files with 115 additions and 52 deletions

View File

@@ -56,7 +56,16 @@ def get_conn():
# ============================================================
def upsert_sync_videos(rows: List[Dict]) -> int:
"""批量 upsert Oracle 传来的 videos 增量。rows 为 Oracle 端 dict 列表。"""
"""批量 upsert Oracle 传来的 videos 增量。rows 为 Oracle 端 dict 列表。
v2 修复 (2026-09-03):不再以 Oracle 的 videos.id 作为 NAS 主键。Oracle 库重建/
重排后 id 会复用/跳变(实测 2303 段跳到 3408/4600+ 段),旧写法按 id 插入时与
已存在的同 filename 行在 UNIQUE filename 上二次冲突报 1062整批中止、游标不推进。
现改为:以 filename 为业务唯一键去重,命中则就地 UPDATE保留 NAS 原 id防止
sync_events/model_calls/identity_map 的 video_id 外键失效Oracle id 仅落
oracle_id 列溯源;未命中则插入(优先用 Oracle id 作主键以对齐子表引用,主键冲突时
回退本地自增,避免 1062
"""
if not rows:
return 0
conn = get_conn()
@@ -64,35 +73,57 @@ def upsert_sync_videos(rows: List[Dict]) -> int:
try:
cur = conn.cursor()
for r in rows:
cur.execute(
"""INSERT INTO sync_videos
(id, drive_file_id, filename, camera_name, duration_sec,
event_start_time, status, summary_json, events_json,
people_json, compute_provider, created_at, updated_at,
processed_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())
ON DUPLICATE KEY UPDATE
drive_file_id=VALUES(drive_file_id),
filename=VALUES(filename),
camera_name=VALUES(camera_name),
duration_sec=VALUES(duration_sec),
event_start_time=VALUES(event_start_time),
status=VALUES(status),
summary_json=VALUES(summary_json),
events_json=VALUES(events_json),
people_json=VALUES(people_json),
compute_provider=VALUES(compute_provider),
created_at=VALUES(created_at),
updated_at=VALUES(updated_at),
processed_at=VALUES(processed_at),
synced_at=NOW()""",
(r.get('id'), r.get('drive_file_id'), r.get('filename'),
r.get('camera_name'), r.get('duration_sec') or 0,
r.get('event_start_time'), r.get('status'),
r.get('summary_json'), r.get('events_json'), r.get('people_json'),
r.get('compute_provider'), r.get('created_at'),
r.get('updated_at'), r.get('processed_at'))
)
fn = r.get('filename')
oracle_id = r.get('id')
cur.execute("SELECT id FROM sync_videos WHERE filename=%s", (fn,))
row = cur.fetchone()
if row:
# 命中业务键:就地更新,保留原 NAS id外键不失效
cur.execute(
"""UPDATE sync_videos SET
oracle_id=%s, drive_file_id=%s, camera_name=%s,
duration_sec=%s, event_start_time=%s, status=%s,
summary_json=%s, events_json=%s, people_json=%s,
compute_provider=%s, created_at=%s, updated_at=%s,
processed_at=%s, synced_at=NOW()
WHERE id=%s""",
(oracle_id, r.get('drive_file_id'), r.get('camera_name'),
r.get('duration_sec') or 0, r.get('event_start_time'),
r.get('status'), r.get('summary_json'), r.get('events_json'),
r.get('people_json'), r.get('compute_provider'),
r.get('created_at'), r.get('updated_at'),
r.get('processed_at'), row[0]))
else:
# 未命中:新文件。优先用 Oracle id 作主键(对齐子表 video_id 引用)
try:
cur.execute(
"""INSERT INTO sync_videos
(id, oracle_id, filename, camera_name, duration_sec,
event_start_time, status, summary_json, events_json,
people_json, compute_provider, created_at, updated_at,
processed_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""",
(oracle_id, oracle_id, fn, r.get('camera_name'),
r.get('duration_sec') or 0, r.get('event_start_time'),
r.get('status'), r.get('summary_json'), r.get('events_json'),
r.get('people_json'), r.get('compute_provider'),
r.get('created_at'), r.get('updated_at'),
r.get('processed_at')))
except Exception:
# 主键 id 被复用(极端情况):回退本地自增,避免 1062
cur.execute(
"""INSERT INTO sync_videos
(oracle_id, filename, camera_name, duration_sec,
event_start_time, status, summary_json, events_json,
people_json, compute_provider, created_at, updated_at,
processed_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""",
(oracle_id, fn, r.get('camera_name'),
r.get('duration_sec') or 0, r.get('event_start_time'),
r.get('status'), r.get('summary_json'), r.get('events_json'),
r.get('people_json'), r.get('compute_provider'),
r.get('created_at'), r.get('updated_at'),
r.get('processed_at')))
n += 1
conn.commit()
return n
@@ -304,7 +335,12 @@ def query_sync_events_for_person_date(person: str, date_str: str) -> List[Dict]:
# ============================================================
def upsert_sync_people(rows: List[Dict]) -> int:
"""批量 upsert Oracle 传来的 people 增量。"""
"""批量 upsert Oracle 传来的 people 增量。
v2 修复 (2026-09-03):同 videos以 label 为业务唯一键去重,命中就地 UPDATE
保留 NAS 原 idpeople.id 无外键引用仅作展示排序Oracle id 落 oracle_id 列
溯源;未命中插入(优先 Oracle id主键冲突回退自增避免 1062
"""
if not rows:
return 0
conn = get_conn()
@@ -312,25 +348,44 @@ def upsert_sync_people(rows: List[Dict]) -> int:
try:
cur = conn.cursor()
for r in rows:
cur.execute(
"""INSERT INTO sync_people
(id, label, canonical_name, first_seen, appearances,
source, features_json, display_uid, updated_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())
ON DUPLICATE KEY UPDATE
label=VALUES(label),
canonical_name=VALUES(canonical_name),
first_seen=VALUES(first_seen),
appearances=VALUES(appearances),
source=VALUES(source),
features_json=VALUES(features_json),
display_uid=VALUES(display_uid),
updated_at=VALUES(updated_at),
synced_at=NOW()""",
(r.get('id'), r.get('label'), r.get('canonical_name'),
r.get('first_seen'), r.get('appearances') or 0,
r.get('source'), r.get('features_json'),
r.get('display_uid'), r.get('updated_at')))
label = r.get('label')
oracle_id = r.get('id')
cur.execute("SELECT id FROM sync_people WHERE label=%s", (label,))
row = cur.fetchone()
if row:
cur.execute(
"""UPDATE sync_people SET
oracle_id=%s, canonical_name=%s, first_seen=%s,
appearances=%s, source=%s, features_json=%s,
display_uid=%s, updated_at=%s, synced_at=NOW()
WHERE id=%s""",
(oracle_id, r.get('canonical_name'), r.get('first_seen'),
r.get('appearances') or 0, r.get('source'),
r.get('features_json'), r.get('display_uid'),
r.get('updated_at'), row[0]))
else:
try:
cur.execute(
"""INSERT INTO sync_people
(id, oracle_id, label, canonical_name, first_seen,
appearances, source, features_json, display_uid,
updated_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""",
(oracle_id, oracle_id, label, r.get('canonical_name'),
r.get('first_seen'), r.get('appearances') or 0,
r.get('source'), r.get('features_json'),
r.get('display_uid'), r.get('updated_at')))
except Exception:
cur.execute(
"""INSERT INTO sync_people
(oracle_id, label, canonical_name, first_seen,
appearances, source, features_json, display_uid,
updated_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""",
(oracle_id, label, r.get('canonical_name'),
r.get('first_seen'), r.get('appearances') or 0,
r.get('source'), r.get('features_json'),
r.get('display_uid'), r.get('updated_at')))
n += 1
conn.commit()
return n