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
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com> Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
113 lines
4.6 KiB
Python
113 lines
4.6 KiB
Python
"""一次性脚本:对历史 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())
|