#!/usr/bin/env python3 """重试 demand_grade_plan_group_item 中 status=failed 的分级任务。 用法: .venv/bin/python scripts/retry_failed_grade_plan_items.py .venv/bin/python scripts/retry_failed_grade_plan_items.py --biz-dt 20260721 .venv/bin/python scripts/retry_failed_grade_plan_items.py --biz-dt 20260721 --workers 5 .venv/bin/python scripts/retry_failed_grade_plan_items.py --biz-dt 20260721 --dry-run .venv/bin/python scripts/retry_failed_grade_plan_items.py --group-id 12 --group-id 15 """ from __future__ import annotations import argparse import json import logging import sys from pathlib import Path _ROOT = Path(__file__).resolve().parents[1] if str(_ROOT) not in sys.path: sys.path.insert(0, str(_ROOT)) from supply_infra.scheduler.jobs.grade_demand_pool import retry_failed_plan_group_items from supply_infra.scheduler.plan_group_batch import MAX_DEMANDS_PER_BATCH logger = logging.getLogger(__name__) def main(argv: list[str] | None = None) -> int: parser = argparse.ArgumentParser( description="重试 demand_grade_plan_group_item 中失败的分级任务", ) parser.add_argument("--biz-dt", default=None, help="业务日期 YYYYMMDD,默认当天") parser.add_argument("--workers", type=int, default=5, help="并发执行的 plan_group 数") parser.add_argument( "--max-demands-per-batch", type=int, default=MAX_DEMANDS_PER_BATCH, help=f"每个 Agent 子批次最多处理的需求条数,默认 {MAX_DEMANDS_PER_BATCH}", ) parser.add_argument( "--group-id", type=int, action="append", dest="group_ids", help="仅重试指定 group_id,可重复传入", ) parser.add_argument( "--dry-run", action="store_true", help="仅列出将要重试的 failed 记录,不实际执行", ) parser.add_argument("--json", action="store_true", help="以 JSON 打印结果") args = parser.parse_args(argv) logging.basicConfig( level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s: %(message)s", ) result = retry_failed_plan_group_items( args.biz_dt, workers=max(1, int(args.workers)), max_demands_per_batch=max(1, min(int(args.max_demands_per_batch), MAX_DEMANDS_PER_BATCH)), group_ids=args.group_ids, dry_run=bool(args.dry_run), ) if args.json: print(json.dumps(result, ensure_ascii=False, indent=2, default=str)) else: reset = result.get("reset") or {} print("\n=== 失败任务重试 ===") print(f"biz_dt={result.get('biz_dt')}") print(f"dry_run={result.get('dry_run')}") print(f"failed_items={result.get('failed_items', 0)}") print(f"reset_items={reset.get('reset_items', 0)}") print(f"reset_groups={reset.get('reset_groups', 0)}") if reset.get("group_ids"): print(f"group_ids={reset.get('group_ids')}") if not result.get("dry_run"): print(f"graded: {result.get('graded_before')} -> {result.get('graded_after')}") print(f"remaining_failed={result.get('remaining_failed', 0)}") print(f"group_status={result.get('group_status')}") print(f"success={result.get('success')}") elif reset.get("items"): print("\n待重试明细:") for item in reset["items"][:20]: print( f" item_id={item['item_id']} group_id={item['group_id']} " f"demand={item['demand_name']!r}" ) if len(reset["items"]) > 20: print(f" ... 另有 {len(reset['items']) - 20} 条") return 0 if result.get("success") else 1 if __name__ == "__main__": raise SystemExit(main())