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