fix(#1714): 同名兜底去重误杀新视频 + ingest 链路孤儿清理/启动恢复 #1734
Reference in New Issue
Block a user
Delete Branch "fix/fallback-dedup-false-positive-1714"
Deleting a branch is permanent. Although the deleted branch may continue to exist for a short time before it actually gets removed, it CANNOT be undone in most cases. Continue?
见 commit message:兜底去重收紧(file_size 严格/hash 非空跳过兜底) + ingest job/asset processing 超时兜底(beat 10min 巡检, processing>60min 转 failed) + worker 启动恢复(CAS 重派卡死 processing, Redis 锁) + celery task_reject_on_worker_lost/visibility_timeout=4h。diff coverage 100%。前端配合:completeDirectUpload 补 file_size=file.size。
根因(staging 实证 IMG_2285.MOV file_size=0 卡 processing 持续误杀): direct complete 不传 file_size,find_recent_active_by_library_and_name 的大小校验 if file_size>0 不生效,30 分钟内同名视频(iPhone IMG_xxxx.MOV) 即使内容全新也被同名兜底误判重复跳过。 兜底去重收紧(宁可漏判不可误杀): - _find_duplicate_asset:file_hash/client_upload_id 非空时不走同名兜底 (hash 已代表内容);走到兜底必须 file_size>0 且与记录大小严格一致 - find_recent_active_by_library_and_name:file_size=0 直接返回 None, 大小条件改为 SQL 内严格等值匹配 - complete 的 _create_pending_asset 补传 file_size(之前占位记录大小永远 0) ingest 链路孤儿兜底(此前只有 generation 链路有清理): - 新增 packages/application/ingest_orphan_cleanup.py: - processing>60min / pending>90min 的 ingest_job 标 failed,关联 processing/uploading asset 联动标 error;无 job 关联超 120min 孤儿 占位 asset 也标 error(beat 每 10 分钟巡检) - worker 启动恢复:processing 超 10 分钟的 job CAS 重置 pending 并 重新派单(Redis SET NX 锁互斥双 worker,旧消息重投由执行前守卫丢弃) - celery 配置:task_reject_on_worker_lost=True; broker visibility_timeout=4h(acks_late 下长转码任务不被误重投) 单测:同名不同大小放行 / hash 非空不走兜底 / file_size=0 放行 / 同名同大小 processing 才判重;孤儿清理 9 例、启动恢复 4 例、 仓储严格守卫 6 例、beat 串联 2 例。diff coverage 100%。 前端配合(前端工程师):completeDirectUpload 请求体补 file_size=file.size。🚀 预览环境已部署
代码审查结果 - PR #1734
⚠️ 问题(1个需要修改)
recover_stuck_ingest_jobs_on_startup函数中,使用了session.query(...).update(...)进行 CAS(Compare-And-Swap)更新。这种方式执行的是原生 SQL UPDATE,不会同步更新 Python 内存中已加载的job_model对象的属性。随后代码修改了job_model.celery_task_id,使得对象变为“脏”状态。在最后的session.commit()时,SQLAlchemy 会检测到job_model的状态变更并执行 Flush,这会发出一条UPDATE ... SET status='processing', celery_task_id=...的语句,从而覆盖了之前 CAS 更新将状态设置为pending的操作。最终导致数据库状态被重置回processing,恢复任务失败。pending并重新派单,导致任务永久卡死。session.refresh(job_model)以同步数据库最新状态到内存对象,然后再修改celery_task_id并提交。💡 建议(1个可选)
make_redis_recovery_lock内部每次调用都创建一个新的redis_lib.Redis.from_url连接。虽然函数有try-except保护且调用频率低(仅启动时一次),但最佳实践是复用连接池或确保连接正确关闭。建议将 Redis 客户端初始化移至模块顶层或使用依赖注入方式。✅ 格式检查通过 | ❌ 逻辑审查需修改 | ✅ 性能无明显问题
🤖 由 AI 代码审查机器人自动生成 | 2026-09-06 05:35:24 | 模型:
🗑️ 预览环境已清理
PR #1734 已关闭或合并,对应的预览环境已被清理。