| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100 |
- #!/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())
|