#!/usr/bin/env python3 """批量执行 demand_grade_plan_group 分级任务。 与定时任务共用 supply_infra.scheduler.jobs.grade_demand_pool.grade_demand_pool。 Usage: .venv/bin/python scripts/run_grade_plan_groups.py .venv/bin/python scripts/run_grade_plan_groups.py --biz-dt 20260721 .venv/bin/python scripts/run_grade_plan_groups.py --biz-dt 20260721 --workers 5 .venv/bin/python scripts/run_grade_plan_groups.py --biz-dt 20260721 --with-orchestrate """ 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.plan_group_batch import MAX_DEMANDS_PER_BATCH from supply_infra.scheduler.jobs.grade_demand_pool import grade_demand_pool logger = logging.getLogger(__name__) def main(argv: list[str] | None = None) -> int: parser = argparse.ArgumentParser(description="批量执行 demand_grade_plan_group 分级任务") parser.add_argument("--biz-dt", default="20260721", help="业务日期 YYYYMMDD,默认 20260721") parser.add_argument( "--max-demands-per-batch", type=int, default=MAX_DEMANDS_PER_BATCH, help=f"每个 Agent 子批次最多处理的需求条数,默认 {MAX_DEMANDS_PER_BATCH}", ) parser.add_argument("--workers", type=int, default=5, help="并发执行的 plan_group 数") parser.add_argument( "--max-rounds", type=int, default=0, help="最多执行轮数,0 表示直到没有 pending 任务", ) parser.add_argument( "--with-orchestrate", action="store_true", help="执行前先跑统筹 Agent 生成/补充计划", ) 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 = grade_demand_pool( str(args.biz_dt).strip(), workers=max(1, int(args.workers)), max_demands_per_batch=max(1, min(int(args.max_demands_per_batch), MAX_DEMANDS_PER_BATCH)), with_orchestrate=bool(args.with_orchestrate), max_rounds=max(0, int(args.max_rounds)), ) if args.json: print(json.dumps(result, ensure_ascii=False, indent=2, default=str)) else: print("\n=== 批量分级完成 ===") print(f"biz_dt={result.get('biz_dt')}") print(f"完成任务组={result.get('groups_run')}") print(f"已分级: {result.get('graded_before')} -> {result.get('graded_after')}") print(f"任务状态: {(result.get('group_status') or result.get('plan_execution', {}).get('final_snapshot', {}).get('group_status'))}") print(f"是否全部完成: {result.get('success')}") return 0 if result.get("success") else 1 if __name__ == "__main__": raise SystemExit(main())