retry_failed_grade_plan_items.py 3.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100
  1. #!/usr/bin/env python3
  2. """重试 demand_grade_plan_group_item 中 status=failed 的分级任务。
  3. 用法:
  4. .venv/bin/python scripts/retry_failed_grade_plan_items.py
  5. .venv/bin/python scripts/retry_failed_grade_plan_items.py --biz-dt 20260721
  6. .venv/bin/python scripts/retry_failed_grade_plan_items.py --biz-dt 20260721 --workers 5
  7. .venv/bin/python scripts/retry_failed_grade_plan_items.py --biz-dt 20260721 --dry-run
  8. .venv/bin/python scripts/retry_failed_grade_plan_items.py --group-id 12 --group-id 15
  9. """
  10. from __future__ import annotations
  11. import argparse
  12. import json
  13. import logging
  14. import sys
  15. from pathlib import Path
  16. _ROOT = Path(__file__).resolve().parents[1]
  17. if str(_ROOT) not in sys.path:
  18. sys.path.insert(0, str(_ROOT))
  19. from supply_infra.scheduler.jobs.grade_demand_pool import retry_failed_plan_group_items
  20. from supply_infra.scheduler.plan_group_batch import MAX_DEMANDS_PER_BATCH
  21. logger = logging.getLogger(__name__)
  22. def main(argv: list[str] | None = None) -> int:
  23. parser = argparse.ArgumentParser(
  24. description="重试 demand_grade_plan_group_item 中失败的分级任务",
  25. )
  26. parser.add_argument("--biz-dt", default=None, help="业务日期 YYYYMMDD,默认当天")
  27. parser.add_argument("--workers", type=int, default=5, help="并发执行的 plan_group 数")
  28. parser.add_argument(
  29. "--max-demands-per-batch",
  30. type=int,
  31. default=MAX_DEMANDS_PER_BATCH,
  32. help=f"每个 Agent 子批次最多处理的需求条数,默认 {MAX_DEMANDS_PER_BATCH}",
  33. )
  34. parser.add_argument(
  35. "--group-id",
  36. type=int,
  37. action="append",
  38. dest="group_ids",
  39. help="仅重试指定 group_id,可重复传入",
  40. )
  41. parser.add_argument(
  42. "--dry-run",
  43. action="store_true",
  44. help="仅列出将要重试的 failed 记录,不实际执行",
  45. )
  46. parser.add_argument("--json", action="store_true", help="以 JSON 打印结果")
  47. args = parser.parse_args(argv)
  48. logging.basicConfig(
  49. level=logging.INFO,
  50. format="%(asctime)s %(levelname)s %(name)s: %(message)s",
  51. )
  52. result = retry_failed_plan_group_items(
  53. args.biz_dt,
  54. workers=max(1, int(args.workers)),
  55. max_demands_per_batch=max(1, min(int(args.max_demands_per_batch), MAX_DEMANDS_PER_BATCH)),
  56. group_ids=args.group_ids,
  57. dry_run=bool(args.dry_run),
  58. )
  59. if args.json:
  60. print(json.dumps(result, ensure_ascii=False, indent=2, default=str))
  61. else:
  62. reset = result.get("reset") or {}
  63. print("\n=== 失败任务重试 ===")
  64. print(f"biz_dt={result.get('biz_dt')}")
  65. print(f"dry_run={result.get('dry_run')}")
  66. print(f"failed_items={result.get('failed_items', 0)}")
  67. print(f"reset_items={reset.get('reset_items', 0)}")
  68. print(f"reset_groups={reset.get('reset_groups', 0)}")
  69. if reset.get("group_ids"):
  70. print(f"group_ids={reset.get('group_ids')}")
  71. if not result.get("dry_run"):
  72. print(f"graded: {result.get('graded_before')} -> {result.get('graded_after')}")
  73. print(f"remaining_failed={result.get('remaining_failed', 0)}")
  74. print(f"group_status={result.get('group_status')}")
  75. print(f"success={result.get('success')}")
  76. elif reset.get("items"):
  77. print("\n待重试明细:")
  78. for item in reset["items"][:20]:
  79. print(
  80. f" item_id={item['item_id']} group_id={item['group_id']} "
  81. f"demand={item['demand_name']!r}"
  82. )
  83. if len(reset["items"]) > 20:
  84. print(f" ... 另有 {len(reset['items']) - 20} 条")
  85. return 0 if result.get("success") else 1
  86. if __name__ == "__main__":
  87. raise SystemExit(main())