Files
xiaoxia-saas/apps/worker/worker_app/tasks/cleanup.py
T
xiaoxia 21c26b5b26
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 2s
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 3s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 3s
CI/CD Pipeline / Check push changed paths (push) Successful in 5s
CI/CD Pipeline / PR Build Web Image (pull_request) Successful in 21s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 23s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 24s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 32s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 32s
CI/CD Pipeline / Build Staging API Image (push) Successful in 32s
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Successful in 7s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 1m54s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m14s
CI/CD Pipeline / Integration Tests (push) Successful in 2m31s
CI/CD Pipeline / Validate - Style (push) Successful in 2m58s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 1m30s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m41s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 1m31s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 1m56s
CI/CD Pipeline / Frontend Unit Tests (push) Successful in 5m29s
CI/CD Pipeline / Validate - Security (push) Successful in 6m18s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m14s
AI Code Review / AI Code Review (pull_request) Successful in 6m28s
CI/CD Pipeline / Unit Tests (push) Successful in 8m25s
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Failing after 267h38m7s
CI/CD Pipeline / CI Gate (push) Failing after 267h38m11s
CI/CD Pipeline / Build Production Worker Image (push) Failing after 267h38m11s
CI/CD Pipeline / Build Production API Image (push) Failing after 267h38m11s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Failing after 267h42m47s
CI/CD Pipeline / Canary Release to Production (pull_request) Failing after 267h45m1s
CI/CD Pipeline / Staging API Integration Tests (pull_request) Failing after 267h45m6s
CI/CD Pipeline / Staging E2E Tests (pull_request) Failing after 267h45m7s
CI/CD Pipeline / Deploy Production (pull_request) Failing after 267h45m10s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Failing after 267h45m10s
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Failing after 267h45m15s
CI/CD Pipeline / Retag skipped Staging Web Image (push) Failing after 267h45m15s
CI/CD Pipeline / Build Production Worker Image (pull_request) Failing after 267h45m20s
CI/CD Pipeline / Retag skipped Staging API Image (push) Failing after 267h45m19s
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Failing after 267h45m26s
CI/CD Pipeline / Build Production Web Image (pull_request) Failing after 267h45m23s
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Failing after 267h45m26s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Failing after 267h46m35s
CI/CD Pipeline / Frontend Lint (pull_request) Failing after 267h46m35s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Failing after 267h46m36s
CI/CD Pipeline / Integration Tests (pull_request) Failing after 267h46m35s
CI/CD Pipeline / Validate - Security (pull_request) Failing after 267h46m36s
CI/CD Pipeline / Frontend Lint (push) Failing after 267h46m38s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 267h46m36s
CI/CD Pipeline / PR Build Worker Image (push) Failing after 267h46m39s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Failing after 267h46m36s
CI/CD Pipeline / Check push changed paths (pull_request) Failing after 267h46m42s
CI/CD Pipeline / Build Staging Web Image (pull_request) Failing after 267h46m37s
CI/CD Pipeline / Build Staging API Image (pull_request) Failing after 267h46m37s
CI/CD Pipeline / PR Build Web Image (push) Failing after 267h46m39s
CI/CD Pipeline / PR Build API Image (push) Failing after 267h46m40s
CI/CD Pipeline / Check if frontend-only change (push) Failing after 267h46m43s
CI/CD Pipeline / Deploy Production (push) Failing after 268h12m38s
CI/CD Pipeline / Build Production Web Image (push) Failing after 268h12m42s
CI/CD Pipeline / ACR Image Cleanup (pull_request) Failing after 268h19m32s
CI/CD Pipeline / Build Production API Image (pull_request) Failing after 268h19m55s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Failing after 268h19m58s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 268h21m6s
feat(generation): worker孤儿任务自动恢复 + 429限流结构化提示 (#1677) (#1710)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-09-05 11:37:49 +08:00

70 lines
2.7 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.
"""定期清理任务 — Celery Beat 调度。
包含:
- cleanup_stale_pending_tasks: 定期清理卡在 pending 超时的 generation_tasks(worker 停止消费时占位)
- cleanup_stale_running_tasks: 定期清理卡在 running 超时的 generation_tasks(容器重启/进程被杀后的孤儿)
"""
import logging
from celery import shared_task
from worker_app.tasks._startup import (
ORPHAN_TASK_TIMEOUT_MINUTES,
PENDING_TASK_TIMEOUT_MINUTES,
cleanup_orphan_tasks,
cleanup_stale_jobs,
cleanup_stale_pending_tasks,
)
logger = logging.getLogger(__name__)
@shared_task(name="worker.cleanup_stale_pending_tasks")
def scheduled_cleanup_stale_pending(timeout_minutes: int = PENDING_TASK_TIMEOUT_MINUTES) -> dict:
"""Celery Beat 调度的定期任务:清理超时的 pending 任务。
每 5 分钟执行一次(由 celery_app.py 的 beat_schedule 配置),
查找所有 status='pending' 且 created_at < NOW() - timeout_minutes
的 generation_tasks,批量更新为 failed,释放限流名额。
Args:
timeout_minutes: 超时时间(分钟),默认 15 分钟
Returns:
{"cleaned": int}
"""
count = cleanup_stale_pending_tasks(timeout_minutes)
if count > 0:
logger.info("[Beat] 清理了 %d 个超时 pending 任务(超时阈值 %d 分钟)", count, timeout_minutes)
return {"cleaned": count}
@shared_task(name="worker.cleanup_stale_running_tasks")
def scheduled_cleanup_stale_running(timeout_minutes: int = ORPHAN_TASK_TIMEOUT_MINUTES) -> dict:
"""Celery Beat 调度的定期任务:清理超时的 running 孤儿任务。
每 5 分钟执行一次。worker_ready 信号只在 worker 启动时清一次,
若 worker 没重启但任务卡死(上传挂起、进程 OOM 被内核杀掉等),
任务会永久卡在 running 占位。此任务做持续兜底:
查找 status='running' 且 updated_at < NOW() - timeout_minutes 的任务,
标记为 failed(原因:容器重启/超时中断),同时清理 Job 表孤儿。
Args:
timeout_minutes: 超时时间(分钟),默认 20 分钟
(worker.generate_video 硬超时 11 分钟,正常任务不可能超过 20 分钟)
Returns:
{"generation_tasks": int, "jobs": int}
"""
gen_count = cleanup_orphan_tasks(timeout_minutes)
job_count = cleanup_stale_jobs(timeout_minutes)
total = gen_count + job_count
if total > 0:
logger.warning(
"[Beat] 清理孤儿任务: running GenerationTask=%d, Job=%d(超时阈值 %d 分钟)",
gen_count,
job_count,
timeout_minutes,
)
return {"generation_tasks": gen_count, "jobs": job_count}