Files
xiaoxia-saas/apps/worker/scripts/backfill_asset_quality.py
T
xiaoxia 01991f14d7
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 3s
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 3s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 3s
CI/CD Pipeline / PR Build API Image (push) Has been skipped
CI/CD Pipeline / PR Build Web Image (push) Has been skipped
CI/CD Pipeline / PR Build Worker Image (push) Has been skipped
CI/CD Pipeline / Check push changed paths (push) Successful in 13s
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / Validate - Style (pull_request) Has been skipped
CI/CD Pipeline / Validate - Security (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been skipped
CI/CD Pipeline / Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 1m3s
CI/CD Pipeline / Build Staging API Image (push) Successful in 1m7s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m9s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 3m13s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m24s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 2m17s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (push) Successful in 3m33s
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been skipped
CI/CD Pipeline / Integration Tests (push) Successful in 4m17s
CI/CD Pipeline / Validate - Style (push) Successful in 4m51s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 4m53s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 55s
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 5m42s
AI Code Review / AI Code Review (pull_request) Successful in 6m57s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 1m52s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m45s
CI/CD Pipeline / PR Build Web Image (pull_request) Failing after 9m5s
CI/CD Pipeline / CI Gate (pull_request) Failing after 2s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 4m58s
CI/CD Pipeline / Unit Tests (push) Successful in 11m33s
CI/CD Pipeline / Validate - Security (push) Successful in 11m53s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / CI Gate (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
fix(queue): #2073 Worker 队列分流——任务路由补全 + beat 独立 + transcode 并发独立 (#2079)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-09-28 01:09:29 +08:00

113 lines
4.6 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.
"""一次性脚本:对历史 quality_score 缺失的视频素材重新打分。
背景(#2073):镜像 97ad0ae2 时期 calculate_quality_score / classify_from_analysis
返回 str 而非 AssetClassification 枚举,导致 calculate_asset_quality 连续报
"'str' object has no attribute 'value'",大量视频素材的 quality_score 卡在 NULL。
镜像 8abdeb95 已修复枚举 bug,但历史失败记录不会自动重跑。本脚本扫描全表,
把 quality_score IS NULL 的视频素材重新投递到 worker.calculate_asset_quality 任务。
使用方式(在 worker 容器内执行):
cd /app/apps/worker
# 干跑,只打印会重跑多少条,不发任务
python -m scripts.backfill_asset_quality --dry-run
# 正式执行
python -m scripts.backfill_asset_quality
# 只重跑最近 N 天的
python -m scripts.backfill_asset_quality --since-days 30
# 限流:每投递一批 sleep 几秒,避免瞬间打爆 transcode 队列
python -m scripts.backfill_asset_quality --batch-size 50 --sleep 2
也可以直接在 staging 机器上 exec 进容器:
docker exec -e PYTHONPATH=/app:/app/apps/api:/app/packages xiaoxia-worker-staging \
python -m scripts.backfill_asset_quality --dry-run
"""
from __future__ import annotations
import argparse
# 保证可以以 python -m scripts.xxx 在容器 /app/apps/worker 下执行
# 也兼容在 repo 根目录下执行(注入路径)
import os
import sys
import time
from datetime import UTC, datetime, timedelta
_SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__))
_WORKER_DIR = os.path.dirname(_SCRIPT_DIR) # apps/worker
_APPS_DIR = os.path.dirname(_WORKER_DIR) # apps
_REPO_ROOT = os.path.dirname(_APPS_DIR) # repo root
for p in (_REPO_ROOT, os.path.join(_REPO_ROOT, "apps", "api"), _REPO_ROOT):
if p not in sys.path:
sys.path.insert(0, p)
def main() -> int:
parser = argparse.ArgumentParser(description="补打历史视频素材 quality_score")
parser.add_argument("--dry-run", action="store_true", help="只统计数量,不投递任务")
parser.add_argument("--since-days", type=int, default=0, help="只处理最近 N 天上传的素材(0=全部)")
parser.add_argument("--batch-size", type=int, default=50, help="每批投递数量,默认 50")
parser.add_argument("--sleep", type=float, default=1.0, help="批次之间 sleep 秒数,默认 1s")
parser.add_argument("--queue", type=str, default="transcode", help="投递队列(默认 transcode)")
args = parser.parse_args()
# 延迟 import,避免在 dry-run 时依赖完整 DB 环境
from worker_app.celery_app import celery_app
from worker_app.db import SessionLocal
from packages.adapters.sqlalchemy_impl.models import AssetModel
db = SessionLocal()
try:
q = db.query(AssetModel).filter(
AssetModel.file_type == "video",
AssetModel.quality_score.is_(None),
)
if args.since_days > 0:
cutoff = datetime.now(UTC) - timedelta(days=args.since_days)
q = q.filter(AssetModel.created_at >= cutoff)
# 先 count 打印
total = q.count()
print(
f"[backfill] 待重跑 quality_score 的视频素材: {total} 条"
f"{' (dry-run,不投递)' if args.dry_run else ''}"
f"{' (最近 ' + str(args.since_days) + ' 天)' if args.since_days > 0 else ''}",
flush=True,
)
if total == 0 or args.dry_run:
return 0
# 分批投递
submitted = 0
batch = 0
offset = 0
while True:
assets = q.order_by(AssetModel.created_at.desc()).offset(offset).limit(args.batch_size).all()
if not assets:
break
batch += 1
for a in assets:
try:
celery_app.send_task(
"worker.calculate_asset_quality",
args=[a.id],
queue=args.queue,
)
submitted += 1
except Exception as e: # noqa: BLE001
print(f"[backfill] 投递失败 asset_id={a.id}: {e}", flush=True)
print(f"[backfill] batch {batch}: 已累计投递 {submitted}/{total}", flush=True)
offset += len(assets)
if args.sleep > 0 and offset < total:
time.sleep(args.sleep)
print(f"[backfill] 完成,共投递 {submitted} 条任务到 {args.queue} 队列", flush=True)
return 0
finally:
db.close()
if __name__ == "__main__":
sys.exit(main())