Files
xiaoxia-saas/apps/worker/worker_app/tasks/_startup.py
T

367 lines
15 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Worker 启动时的初始化任务 — 孤儿任务清理等."""
import logging
from datetime import UTC
from celery.signals import worker_ready
from worker_app.db import SessionLocal
logger = logging.getLogger(__name__)
def cleanup_stale_running_with_session(repo, timeout_minutes: int) -> int:
"""清理超时未更新的 running GenerationTask(可注入 repo 的纯核心,便于单测)。
Returns:
清理的任务数量
"""
return len(cleanup_stale_running_with_session_ids(repo, timeout_minutes))
def cleanup_stale_running_with_session_ids(repo, timeout_minutes: int) -> list[tuple[str, str]]:
"""同 cleanup_stale_running_with_session,返回 [(task_id, celery_task_id), ...]。"""
fn = getattr(repo, "cleanup_stale_running_with_ids", None)
if fn is not None:
return fn(timeout_minutes)
# 旧仓储无 _with_ids 方法:降级为计数,无法撤销消息(执行前状态守卫兜底)
count = repo.cleanup_stale_running(timeout_minutes)
return [("", "") for _ in range(count)]
def cleanup_stale_pending_with_session(repo, timeout_minutes: int) -> int:
"""清理超时 pending GenerationTask(可注入 repo 的纯核心,便于单测)。
Returns:
清理的任务数量
"""
return len(cleanup_stale_pending_with_session_ids(repo, timeout_minutes))
def cleanup_stale_pending_with_session_ids(repo, timeout_minutes: int) -> list[tuple[str, str]]:
"""同 cleanup_stale_pending_with_session,返回 [(task_id, celery_task_id), ...]。"""
fn = getattr(repo, "cleanup_stale_pending_with_ids", None)
if fn is not None:
return fn(timeout_minutes)
count = repo.cleanup_stale_pending(timeout_minutes)
return [("", "") for _ in range(count)]
def _revoke_and_purge_stale_messages(items: list[tuple[str, str]]) -> int:
"""把清理掉的任务对应的 Celery 消息撤销并从 Redis 队列清除(#1714)。
防止「DB 已标 failed,但队列消息还在 → 重投执行 → 非法状态转换 → 半成品」。
失败不阻断清理流程(执行前状态守卫是第二道防线)。
"""
biz_ids = [tid for tid, _ in items if tid]
celery_ids = [cid for _, cid in items if cid]
if not biz_ids and not celery_ids:
return 0
try:
from worker_app.celery_app import celery_app as app
from worker_app.core.config import get_settings
from packages.shared.celery_orphan_guard import revoke_and_purge
broker_url = get_settings().broker_url
return revoke_and_purge(
app,
broker_url,
business_task_ids=biz_ids,
celery_task_ids=celery_ids,
)
except Exception as e: # noqa: BLE001
logger.error("撤销作废任务队列消息失败(执行前守卫仍会兜底): %s", e, exc_info=True)
return 0
# 孤儿任务超时阈值:running 任务超过此时间无进度更新则视为卡死。
# 依据:worker.generate_video 硬超时 time_limit=11 分钟,正常任务不可能超过;
# 20 分钟阈值覆盖硬超时 + 重试 + 余量,绝不误杀正常任务。
ORPHAN_TASK_TIMEOUT_MINUTES = 20
# Pending 任务超时阈值:任务创建后超过此时间仍未开始执行则判死。
# 注意区分 running 孤儿阈值(20 分钟):pending 是「排队等待」时间,
# 队列积压(如 20+ 转码任务)时视频生成可能正常排队较久,阈值必须放宽,
# 避免正常排队任务被误杀。队列隔离(#1714)后 generation 队列独占 worker,
# 理论上排队极短;保留 45 分钟作为兜底,覆盖 worker 短暂停止消费的场景。
PENDING_TASK_TIMEOUT_MINUTES = 45
def cleanup_orphan_tasks(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> int: # pragma: no cover
"""清理数据库中超时未更新的 running GenerationTask(孤儿任务)。
worker 重启或崩溃后,之前处于 running 状态的任务会变成孤儿任务,
一直卡在 running 不动。通过 updated_at 超时判断并标记为 failed。
Args:
timeout_minutes: 超时时间(分钟),默认 10 分钟
Returns:
清理的任务数量
"""
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
SQLAlchemyGenerationTaskRepository,
)
try:
session = SessionLocal()
try:
repo = SQLAlchemyGenerationTaskRepository(session)
items = cleanup_stale_running_with_session_ids(repo, timeout_minutes)
finally:
session.close()
count = len(items)
if count > 0:
logger.warning("清理了 %d 个超时的孤儿 GenerationTask(超过 %d 分钟未更新)", count, timeout_minutes)
purged = _revoke_and_purge_stale_messages(items)
logger.info("孤儿任务对应队列消息撤销/清除完成: %d 条", purged)
else:
logger.info("无孤儿 GenerationTask 需要清理")
return count
except Exception as e:
logger.error("清理孤儿 GenerationTask 失败: %s", e, exc_info=True)
return 0
def cleanup_stale_jobs(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> int: # pragma: no cover
"""清理数据库中超时未更新的 running Job(孤儿任务)。
与 cleanup_orphan_tasks 配合,同时清理 Job 表和 GenerationTask 表。
Returns:
清理的任务数量
"""
from datetime import datetime, timedelta
from packages.adapters.sqlalchemy_impl.models import JobModel
from packages.domain.job import JobStatus
try:
session = SessionLocal()
cutoff = datetime.now(UTC) - timedelta(minutes=timeout_minutes)
stale_jobs = (
session.query(JobModel)
.filter(
JobModel.status == JobStatus.RUNNING.value,
JobModel.updated_at < cutoff,
)
.all()
)
count = 0
stale_items: list[tuple[str, str]] = []
for model in stale_jobs:
stale_items.append((model.id, getattr(model, "celery_task_id", "") or ""))
model.status = JobStatus.FAILED.value
model.error_message = f"任务执行中断(超过 {timeout_minutes} 分钟未更新)"
count += 1
if count > 0:
session.commit()
logger.warning("清理了 %d 个超时的孤儿 Job(超过 %d 分钟未更新)", count, timeout_minutes)
else:
logger.info("无孤儿 Job 需要清理")
session.close()
if count > 0:
_revoke_and_purge_generation(stale_items)
return count
except Exception as e:
logger.error("清理孤儿 Job 失败: %s", e, exc_info=True)
return 0
def cleanup_stale_pending_tasks(timeout_minutes: int = PENDING_TASK_TIMEOUT_MINUTES) -> int: # pragma: no cover
"""清理数据库中卡在 pending 状态超时的 GenerationTask。
全局任务队列有 pending 数量上限,长期卡在 pending 的任务会占满队列,
导致新用户无法创建任务。通过 created_at 超时判断并标记为 failed。
Args:
timeout_minutes: 超时时间(分钟),默认 PENDING_TASK_TIMEOUT_MINUTES
Returns:
清理的任务数量
"""
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
SQLAlchemyGenerationTaskRepository,
)
session = SessionLocal()
try:
repo = SQLAlchemyGenerationTaskRepository(session)
items = cleanup_stale_pending_with_session_ids(repo, timeout_minutes)
count = len(items)
if count > 0:
logger.warning("清理了 %d 个超时的 pending GenerationTask(超过 %d 分钟未处理)", count, timeout_minutes)
purged = _revoke_and_purge_stale_messages(items)
logger.info("超时 pending 任务对应队列消息撤销/清除完成: %d 条", purged)
else:
logger.info("无超时 pending GenerationTask 需要清理")
return count
except Exception as e:
logger.error("清理超时 pending GenerationTask 失败: %s", e, exc_info=True)
return 0
finally:
session.close()
def _revoke_and_purge_generation(items: list[tuple[str, str]]) -> int:
"""撤销 Job 表孤儿任务(TTS/配音等)的队列消息,队列覆盖全部已知队列。"""
biz_ids = [tid for tid, _ in items if tid]
celery_ids = [cid for _, cid in items if cid]
if not biz_ids and not celery_ids:
return 0
try:
from worker_app.celery_app import celery_app as app
from worker_app.core.config import get_settings
from packages.shared.celery_orphan_guard import revoke_and_purge
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=("generation", "transcode", "celery"),
)
except Exception as e: # noqa: BLE001
logger.error("撤销 Job 队列消息失败: %s", e, exc_info=True)
return 0
def cleanup_all_stale_tasks(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> dict: # pragma: no cover
"""统一清理所有超时的孤儿任务。
同时清理 GenerationTask 和 Job 两类表。
Returns:
{"generation_tasks": int, "jobs": int}
"""
gen_count = cleanup_orphan_tasks(timeout_minutes)
job_count = cleanup_stale_jobs(timeout_minutes)
pending_count = cleanup_stale_pending_tasks(PENDING_TASK_TIMEOUT_MINUTES)
total = gen_count + job_count + pending_count
if total > 0:
logger.warning(
"任务清理完成: 孤儿 GenerationTask=%d, 孤儿 Job=%d, 超时 pending=%d, 总计=%d",
gen_count,
job_count,
pending_count,
total,
)
return {"generation_tasks": gen_count, "jobs": job_count, "pending": pending_count}
@worker_ready.connect
def _on_worker_ready(sender, **kwargs): # pragma: no cover
"""Worker 启动完成后执行 — 清理孤儿任务。"""
logger.info("Worker 启动完成,开始清理孤儿 running 任务...")
result = cleanup_all_stale_tasks()
total = result["generation_tasks"] + result["jobs"]
logger.info("Worker 启动清理完成,共清理 %d 个孤儿任务", total)
@worker_ready.connect
def _recover_stuck_ingest_jobs_on_ready(sender, **kwargs): # pragma: no cover
"""Worker 启动完成后恢复卡死在 processing 的 ingest_job(#1714)。
容器重启/进程 OOM 导致 transcode 队列 unacked 消息未重投时,processing
ingest_job 会永久卡死。启动时扫描 processing 超 10 分钟的 job,CAS 重置
pending 并重新派单;Redis 锁保证同容器 generation/transcode 双 worker
只有一个执行恢复。旧消息若后来重投,ingest_asset 执行前守卫会丢弃。
"""
try:
from packages.application.ingest_orphan_cleanup import (
make_redis_recovery_lock,
recover_stuck_ingest_jobs_on_startup,
)
session = SessionLocal()
try:
recovered = recover_stuck_ingest_jobs_on_startup(
session,
lock_acquire=make_redis_recovery_lock(),
stuck_minutes=10,
)
finally:
session.close()
logger.info("Worker 启动 ingest 恢复完成,共重新派单 %d 个卡死任务", recovered)
except Exception as e: # noqa: BLE001 — 启动恢复失败不能阻断 worker 起服
logger.error("启动 ingest 恢复扫描失败(beat 巡检仍会兜底标 failed): %s", e, exc_info=True)
def recover_stale_voice_clones_on_startup(timeout_minutes: int = 10) -> int:
"""Worker 启动时恢复卡死在 processing 的音色克隆任务。
容器重启/进程 OOM 时 worker 中正在轮询的克隆任务会丢失,
voice_clone_profiles 永久卡在 processing 无兜底。启动时扫描
updated_at 超过 timeout_minutes 的 processing 记录,直接标记
为 failed(错误信息指引用户重试)。选择标 failed 而非重新派单,
因为 CosyVoice 侧的 voice_id 无法在无上下文下恢复轮询,重试需
用户确认后显式触发。
Args:
timeout_minutes: 判定卡死的阈值,默认 10 分钟
Returns:
恢复的记录数
"""
from packages.adapters.sqlalchemy_impl.voice_clone_profile_repository import (
SQLAlchemyVoiceCloneProfileRepository,
)
try:
session = SessionLocal()
try:
repo = SQLAlchemyVoiceCloneProfileRepository(session)
count = repo.cleanup_stale_processing(timeout_minutes)
finally:
session.close()
if count > 0:
logger.warning("启动时恢复了 %d 个卡死在 processing 的音色克隆(超时 %d 分钟)", count, timeout_minutes)
else:
logger.info("无卡死 processing 音色克隆需要恢复")
return count
except Exception as e:
logger.error("启动时音色克隆恢复扫描失败(beat 巡检仍会兜底): %s", e, exc_info=True)
return 0
@worker_ready.connect
def _recover_stuck_voice_clones_on_ready(sender, **kwargs):
"""Worker 启动完成后恢复卡死的音色克隆任务。"""
try:
recovered = recover_stale_voice_clones_on_startup()
logger.info("Worker 启动音色克隆恢复完成,共标记 %d 个卡死任务为 failed", recovered)
except Exception as e:
logger.error("启动音色克隆恢复失败(beat 巡检仍会兜底标 failed): %s", e, exc_info=True)
@worker_ready.connect
def _probe_gpu_encoder_on_ready(sender, **kwargs):
"""Worker 启动完成后探测 P4000 GPU NVENC 节点状态,打日志。"""
try:
from packages.shared.gpu_encoder import get_gpu_encoder
client = get_gpu_encoder()
if client is None:
logger.info(
"[gpu-encoder] disabled (ENABLE_GPU_ENCODE=false or endpoint not configured), using CPU libx264"
)
return
health = client.check_health()
if health.ready:
logger.info(
"[gpu-encoder] NVENC enabled: endpoint=%s gpu=%s worker=%s",
client.endpoint,
health.gpu_name,
health.worker,
)
else:
logger.warning(
"[gpu-encoder] configured but NOT ready: %s (endpoint=%s) — falling back to CPU",
health.error,
client.endpoint,
)
except Exception as e: # noqa: BLE001
logger.warning("[gpu-encoder] startup probe error (will retry on first job, CPU fallback): %s", e)