Files
xiaoxia-saas/apps/worker/worker_app/tasks/_startup.py
T
CI Bot ddf861b02f
CI Build & Deploy Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI Build & Deploy Pipeline / Build Staging API Image (pull_request) Has been skipped
CI Build & Deploy Pipeline / Build Production API Image (pull_request) Has been skipped
CI Build & Deploy Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI Build & Deploy Pipeline / Build Production Web Image (pull_request) Has been skipped
CI Build & Deploy Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI Build & Deploy Pipeline / Deploy Production (pull_request) Has been skipped
CI Build & Deploy Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI Build & Deploy Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI Build & Deploy Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI Build & Deploy Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 24s
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 55s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 56s
CI/CD Pipeline / Validate Code Quality And Tests (pull_request) Failing after 1m29s
Auto Merge CI PRs / Auto Merge on CI Green + Approved (pull_request) Successful in 2m5s
AI Code Review / AI Code Review (pull_request) Successful in 2m32s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m36s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 3m8s
Auto Approve CI PRs / Auto Approve on CI Green (pull_request) Successful in 3m42s
feat(P0): 孤儿任务清理 - worker启动+任务调度时自动清理超时running任务
- GenerationTaskModel 新增 updated_at 字段(onupdate自动刷新)
- 领域模型 GenerationTask 新增 updated_at
- repo 新增 cleanup_stale_running 方法:running且updated_at超10分钟→failed
- worker启动时(on_worker_ready signal)自动清理
- 任务调度时(generate_video入口)也清理一次
- error_message: 任务执行中断(worker重启/超时)
- 配套5个单元测试
2026-07-18 19:19:09 +08:00

56 lines
1.8 KiB
Python
Executable File

"""Worker 启动时的初始化任务 — 孤儿任务清理等。"""
import logging
from worker_app.celery_app import celery_app
from worker_app.db import SessionLocal
logger = logging.getLogger(__name__)
ORPHAN_TASK_TIMEOUT_MINUTES = 10
def cleanup_orphan_tasks(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> int:
"""清理数据库中超时未更新的 running 任务。
worker 重启或崩溃后,之前处于 running 状态的任务会变成孤儿任务,
一直卡在 running 不动。通过 updated_at 超时判断并标记为 failed。
Args:
timeout_minutes: 超时时间(分钟),默认 10 分钟
Returns:
清理的任务数量
"""
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
SQLAlchemyGenerationTaskRepository,
)
try:
session = SessionLocal()
repo = SQLAlchemyGenerationTaskRepository(session)
count = repo.cleanup_stale_running(timeout_minutes)
session.close()
if count > 0:
logger.warning("清理了 %d 个超时的孤儿 running 任务", count)
else:
logger.info("无孤儿 running 任务需要清理")
return count
except Exception as e:
logger.error("清理孤儿任务失败: %s", e, exc_info=True)
return 0
@celery_app.on_after_configure.connect
def _setup_periodic_tasks(sender, **kwargs):
"""Celery 配置完成后,注册 worker 启动钩子。"""
pass
@celery_app.on_worker_ready.connect
def _on_worker_ready(sender, **kwargs):
"""Worker 启动完成后执行 — 清理孤儿任务。"""
logger.info("Worker 启动完成,开始清理孤儿 running 任务...")
count = cleanup_orphan_tasks()
logger.info("Worker 启动清理完成,共清理 %d 个孤儿任务", count)