-
Notifications
You must be signed in to change notification settings - Fork 490
Expand file tree
/
Copy pathbatch_local_artifacts.py
More file actions
132 lines (113 loc) · 3.94 KB
/
Copy pathbatch_local_artifacts.py
File metadata and controls
132 lines (113 loc) · 3.94 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
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
"""Run a tiny observable micro-batch ETL that writes local artifacts."""
from __future__ import annotations
import argparse
import asyncio
import json
import math
from pathlib import Path
from quantmind.etl import BatchETLPipeline, PipelineContext
BATCH_SIZE = 2
async def extract(source: list[str], *, ctx: PipelineContext):
"""Yield small in-memory batches and report each batch before yielding."""
for offset in range(0, len(source), BATCH_SIZE):
batch = source[offset : offset + BATCH_SIZE]
await ctx.progress(
len(batch),
total=len(batch),
message=f"Prepared rows {offset + 1}-{offset + len(batch)}",
)
yield batch
async def transform(
batch: list[str], *, ctx: PipelineContext
) -> list[dict[str, object]]:
"""Build deterministic records for one batch."""
await ctx.progress(1, total=1, message="Normalized batch")
return [
{
"text": line.strip(),
"characters": len(line.strip()),
}
for line in batch
if line.strip()
]
async def run_example(*, dry_run: bool) -> dict[str, object]:
"""Create the run, print its receipt, then execute all batches."""
artifact_dir = Path.cwd() / ".quant-mind" / "etl-batch-example"
async def load(
records: list[dict[str, object]], *, ctx: PipelineContext
) -> dict[str, int]:
batch_index = ctx.batch_index
if batch_index is None:
raise RuntimeError("batch load requires a batch index")
target = artifact_dir / f"batch-{batch_index:03d}.json"
content = json.dumps(records, indent=2, sort_keys=True) + "\n"
byte_count = len(content.encode("utf-8"))
if ctx.dry_run:
await ctx.progress(
1,
total=1,
message="Planned batch artifact",
metrics={
"planned_artifacts": 1,
"planned_bytes": byte_count,
"planned_records": len(records),
},
)
return {
"planned_artifacts": 1,
"planned_records": len(records),
}
artifact_dir.mkdir(parents=True, exist_ok=True)
temporary = target.with_suffix(".json.tmp")
await asyncio.to_thread(temporary.write_text, content, encoding="utf-8")
temporary.replace(target)
await ctx.progress(
1,
total=1,
message="Wrote batch artifact",
metrics={
"artifacts_written": 1,
"bytes_written": byte_count,
"records_written": len(records),
},
)
return {"artifacts_written": 1, "records_written": len(records)}
source = ["alpha", " beta ", "", "gamma"]
pipeline = BatchETLPipeline(
"batch-local-artifacts",
extract=extract,
transform=transform,
load=load,
)
run = pipeline.create_run(
source,
dry_run=dry_run,
config_summary={
"artifact_format": "json",
"batch_size": BATCH_SIZE,
},
total_batches=math.ceil(len(source) / BATCH_SIZE),
)
print(run.receipt(), flush=True)
summary = await run.execute()
return {
"dry_run": dry_run,
"completed_batches": summary.completed_batches,
"counts": dict(summary.counts),
}
def _parse_args() -> argparse.Namespace:
"""Parse the runtime run-mode option."""
parser = argparse.ArgumentParser()
parser.add_argument(
"--dry-run",
action="store_true",
help="plan and validate without creating example batch artifacts",
)
return parser.parse_args()
async def main() -> None:
"""Run the example in normal or dry-run mode."""
args = _parse_args()
summary = await run_example(dry_run=bool(args.dry_run))
print("summary=" + json.dumps(summary, sort_keys=True))
if __name__ == "__main__":
asyncio.run(main())