"""上传/转码链路(IngestJob + Asset)孤儿清理核心逻辑。 #1714:generation 链路有 cleanup_stale_running/pending 兜底,但上传链路 (ingest_jobs + assets)没有。worker 容器重启/进程 OOM 时,已 prefetch 的 celery 消息会丢失(transcode 队列 worker_prefetch_multiplier=1,消息预取后 宕机即丢失,Redis 队列里也不再存在),导致: - ingest_jobs.status 永久卡 pending/processing - assets.status 永久卡 processing/uploading(complete 阶段预建的占位) 本模块提供纯核心(session 注入,便于单测):超时阈值内无更新的记录 批量标终态(job→failed、asset→error),并返回 (job_id, celery_task_id) 列表供调用方 revoke + purge 残留队列消息。 """ from __future__ import annotations import logging from collections.abc import Callable from datetime import UTC, datetime, timedelta from typing import Any logger = logging.getLogger(__name__) # ingest_job PROCESSING 超时阈值:ingest 任务包含下载 + ffprobe + HEVC 转码 # (1GB 视频约 10-20 分钟)+ 回传 OSS,正常任务可能跑 20-30 分钟; # 60 分钟阈值覆盖大文件转码 + 抖动,绝不误杀正常任务。 INGEST_PROCESSING_TIMEOUT_MINUTES = 60 # ingest_job PENDING 超时阈值:transcode 队列 concurrency=1,队列积压时 # 正常排队可能较久;90 分钟覆盖 worker 短暂停消费 + 排队。 INGEST_PENDING_TIMEOUT_MINUTES = 90 # Asset 占位超时阈值:无关联 ingest_job 的孤儿占位(complete 预建后派单失败等), # 阈值放宽到 120 分钟,避免与 ingest_job 生命周期错杀。 ASSET_ORPHAN_TIMEOUT_MINUTES = 120 _TERMINAL_JOB_STATUSES = ("failed", "completed") _TERMINAL_ASSET_STATUSES = ("ready", "error", "deleted") def _now() -> datetime: return datetime.now(UTC) def cleanup_stale_ingest_jobs( session: Any, *, processing_timeout_minutes: int = INGEST_PROCESSING_TIMEOUT_MINUTES, pending_timeout_minutes: int = INGEST_PENDING_TIMEOUT_MINUTES, commit: bool = True, ) -> tuple[list[tuple[str, str]], list[str]]: """清理超时卡 pending/processing 的 ingest_jobs,并联动关联 asset。 Args: session: SQLAlchemy session(或提供 query/commit 的鸭子类型) processing_timeout_minutes: processing 状态超时阈值 pending_timeout_minutes: pending 状态超时阈值 commit: 是否提交事务 Returns: (job_items, asset_ids) - job_items: [(job_id, celery_task_id), ...] 供 revoke/purge - asset_ids: 被联动标记为 error 的 asset id 列表 """ from packages.adapters.sqlalchemy_impl.models import AssetModel, IngestJobModel now = _now() processing_cutoff = now - timedelta(minutes=processing_timeout_minutes) pending_cutoff = now - timedelta(minutes=pending_timeout_minutes) stale_jobs = ( session.query(IngestJobModel) .filter( IngestJobModel.status.in_(["pending", "processing"]), ( (IngestJobModel.status == "processing") & (IngestJobModel.updated_at < processing_cutoff) | (IngestJobModel.status == "pending") & (IngestJobModel.created_at < pending_cutoff) ), ) .all() ) job_items: list[tuple[str, str]] = [] asset_ids: list[str] = [] stale_asset_models: list[Any] = [] for job_model in stale_jobs: ref_time = job_model.updated_at or job_model.created_at if ref_time.tzinfo is None: # SQLite 读回 naive datetime 的防御 ref_time = ref_time.replace(tzinfo=UTC) stale_minutes = int((now - ref_time).total_seconds() // 60) job_model.status = "failed" job_model.error_message = ( f"转码任务执行中断(超过超时阈值未更新,疑似 worker 重启/进程退出,已卡死 {stale_minutes} 分钟)" ) job_model.updated_at = now job_items.append((job_model.id, getattr(job_model, "celery_task_id", "") or "")) if job_model.asset_id: asset_ids.append(job_model.asset_id) if asset_ids: stale_asset_models = ( session.query(AssetModel) .filter( AssetModel.id.in_(asset_ids), AssetModel.status.in_(["processing", "uploading"]), ) .all() ) for asset_model in stale_asset_models: asset_model.status = "error" asset_model.updated_at = now if commit and (job_items or stale_asset_models): session.commit() if job_items: logger.warning( "[ingest-cleanup] 清理 %d 个超时 ingest_job(processing>%dm / pending>%dm),联动 %d 个 asset 标 error", len(job_items), processing_timeout_minutes, pending_timeout_minutes, len(stale_asset_models), ) return job_items, [a.id for a in stale_asset_models] def cleanup_orphan_processing_assets( session: Any, *, timeout_minutes: int = ASSET_ORPHAN_TIMEOUT_MINUTES, commit: bool = True, ) -> list[str]: """清理无 ingest_job 关联、超时卡 processing/uploading 的孤儿 asset 占位。 complete 阶段预建 asset 后若派单失败(或 direct 上传 complete 后 未触发 ingest),占位会永久卡住。这类 asset 没有对应 ingest_job, 只能按 created_at 超时兜底标 error。 """ from packages.adapters.sqlalchemy_impl.models import AssetModel, IngestJobModel cutoff = _now() - timedelta(minutes=timeout_minutes) orphan_assets = ( session.query(AssetModel) .outerjoin(IngestJobModel, IngestJobModel.asset_id == AssetModel.id) .filter( AssetModel.status.in_(["processing", "uploading"]), AssetModel.created_at < cutoff, IngestJobModel.id.is_(None), ) .all() ) for asset_model in orphan_assets: asset_model.status = "error" asset_model.updated_at = _now() if commit and orphan_assets: session.commit() logger.warning("[ingest-cleanup] 清理 %d 个无 job 关联的超时孤儿 asset 占位", len(orphan_assets)) return [a.id for a in orphan_assets] def revoke_stale_ingest_messages( job_items: list[tuple[str, str]], *, celery_app_factory: Callable[[], Any] | None = None, broker_url_factory: Callable[[], str] | None = None, ) -> int: """revoke + 物理清理 ingest 作废消息(transcode/celery 队列)。 消息可能已在 worker 宕机时丢失(队列里查不到),那也无害; 若消息还在(极端重复投递),物理清除防止重投执行。 失败不阻断清理(ingest_asset 的执行前状态守卫是第二道防线)。 """ biz_ids = [jid for jid, _ in job_items if jid] celery_ids = [cid for _, cid in job_items if cid] if not biz_ids and not celery_ids: return 0 try: from packages.shared.celery_orphan_guard import revoke_and_purge app = celery_app_factory() if celery_app_factory else None broker_url = broker_url_factory() if broker_url_factory else "" if app is None or not broker_url: from worker_app.celery_app import celery_app as _app from worker_app.core.config import get_settings app = _app broker_url = get_settings().broker_url return revoke_and_purge( app, broker_url, business_task_ids=biz_ids, celery_task_ids=celery_ids, queue_names=("transcode", "celery"), ) except Exception as e: # noqa: BLE001 logger.error("撤销作废 ingest 队列消息失败(执行前守卫仍会兜底): %s", e, exc_info=True) return 0 # ── worker 启动恢复(#1714)────────────────────────────────────────────── # # task_acks_late=True 下,worker 崩溃/容器重启时未 ack 的消息理论上会在 # visibility_timeout 到期后重新投递;但 prefork 进程异常、部署窗口跨 # visibility 配置边界等场景仍可能留下卡在 processing 的 ingest_job # (staging 实证:03:16 派单、03:45 置 processing 后 worker 重启, # unacked 消息未重投,任务永久卡死)。启动时做一次显式恢复扫描兜底。 # # 恢复策略:processing 超过 stuck_minutes(默认 10 分钟,部署中跨进程 # 交接的正常窗口 < 10 分钟,不会误抢别的 worker 正在执行的任务)的 job, # CAS 重置为 pending 并重新 send_task;旧消息若后来重投,ingest_asset # 的执行前守卫会把状态不匹配的旧 celery 消息丢弃。 def recover_stuck_ingest_jobs_on_startup( session: Any, *, send_task: Callable[..., Any] | None = None, update_celery_task_id: Callable[[str, str], None] | None = None, lock_acquire: Callable[[], bool] | None = None, stuck_minutes: int = 10, commit: bool = True, ) -> int: """worker 启动时把卡在 processing 超时的 ingest_job 重新派单。 Args: session: SQLAlchemy session send_task: celery send_task 可调用(注入便于测试);不传则用 worker celery_app update_celery_task_id: 回写新 celery task id 的回调(job_id, new_task_id) lock_acquire: 分布式锁获取回调(多 worker 进程同时启动时只允许一个恢复); 返回 False 表示未抢到锁,本次跳过 stuck_minutes: processing 超过该分钟数视为卡死 Returns: 重新派单的 job 数 """ if lock_acquire is not None and not lock_acquire(): logger.info("[ingest-recover] 未抢到恢复锁,跳过(另一进程正在恢复)") return 0 from packages.adapters.sqlalchemy_impl.models import IngestJobModel cutoff = _now() - timedelta(minutes=stuck_minutes) stuck_jobs = ( session.query(IngestJobModel) .filter(IngestJobModel.status == "processing", IngestJobModel.updated_at < cutoff) .order_by(IngestJobModel.updated_at.asc()) .all() ) if not stuck_jobs: logger.info("[ingest-recover] 无卡死 processing ingest_job 需要恢复") return 0 if send_task is None: from worker_app.celery_app import celery_app as _app send_task = _app.send_task recovered = 0 for job_model in stuck_jobs: # CAS:只有仍是 processing 才重置(并发/旧消息已回写终态时不碰) updated = ( session.query(IngestJobModel) .filter(IngestJobModel.id == job_model.id, IngestJobModel.status == "processing") .update({"status": "pending", "error_message": "", "updated_at": _now()}) ) if not updated: continue try: result = send_task("worker.ingest_asset", args=[job_model.id]) new_task_id = getattr(result, "id", "") or "" except Exception as e: # noqa: BLE001 logger.error("[ingest-recover] 重新派单失败 job_id=%s: %s", job_model.id, e) continue if new_task_id: job_model.celery_task_id = new_task_id if update_celery_task_id is not None: update_celery_task_id(job_model.id, new_task_id) logger.warning( "[ingest-recover] 卡死 ingest_job %s 已重置 pending 并重新派单 (new celery task=%s)", job_model.id, new_task_id, ) recovered += 1 if commit and recovered: session.commit() logger.warning("[ingest-recover] 启动恢复完成,共重新派单 %d 个卡死 ingest_job", recovered) return recovered def make_redis_recovery_lock(lock_key: str = "ingest:recover:startup", ttl_seconds: int = 300): """构造基于 Redis SET NX 的恢复锁工厂(多 worker 进程互斥)。 返回一个无参 callable,调用时尝试抢锁:抢到返回 True,未抢到返回 False。 Redis 不可用时不阻断启动恢复(返回 True,恢复逻辑自身有 CAS 幂等保护)。 """ def _acquire() -> bool: try: import redis as redis_lib from worker_app.core.config import get_settings client = redis_lib.Redis.from_url(get_settings().broker_url) return bool(client.set(lock_key, "1", nx=True, ex=ttl_seconds)) except Exception as e: # noqa: BLE001 logger.warning("[ingest-recover] Redis 锁不可用,降级为无锁执行(CAS 兜底): %s", e) return True return _acquire