-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathdistributed_execution_cli.py
More file actions
29 lines (22 loc) · 1.09 KB
/
Copy pathdistributed_execution_cli.py
File metadata and controls
29 lines (22 loc) · 1.09 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
"""CLI demo for QueueCraft distributed execution."""
from __future__ import annotations
import argparse
import json
from pathlib import Path
from distributed_execution import CheckpointStore, DistributedExecutor, DistributedPlan
def main() -> None:
parser = argparse.ArgumentParser()
parser.add_argument("--tasks", type=int, default=100)
parser.add_argument("--workers", type=int, default=4)
parser.add_argument("--chunk-size", type=int, default=10)
parser.add_argument("--seed", type=int, default=42)
parser.add_argument("--checkpoint", default="artifacts/distributed-checkpoint.json")
args = parser.parse_args()
plan = DistributedPlan("CLI-DEMO", args.tasks, args.workers, args.chunk_size, args.seed)
events = []
executor = DistributedExecutor(plan, CheckpointStore(Path(args.checkpoint)))
result = executor.run(lambda task_id, task_seed: {"task_id": task_id, "seed": task_seed}, progress=events.append)
result["last_progress"] = events[-1] if events else None
print(json.dumps(result, ensure_ascii=False, indent=2))
if __name__ == "__main__":
main()