-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathstreaming_handler.py.tmpl
More file actions
52 lines (44 loc) · 2.21 KB
/
Copy pathstreaming_handler.py.tmpl
File metadata and controls
52 lines (44 loc) · 2.21 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
"""Streaming handler example for newline-delimited JSON input.
Copy to raincloud/pipeline/handlers/<name>.py, rename the function, and declare
it in HANDLERS in raincloud/_registry.py ("<name>": "<name>:<name>"), the only
registration. Set parse.reader="custom" in the manifest so parsed contains
paths. Adapt the ingestion query for the upstream shape.
The handler writes canonical IPC in batches and returns []; the builder then
validates and exports that canonical artifact. Run it only through
build.run_one: workdir_root() is recipe-scoped only inside spec.recipe_scratch(),
and open_canonical_writer registers its output with lifecycle.build_outputs,
which publishes it atomically when the build succeeds. run_one sets up both,
holds the operation lock, and writes the build record.
For VARIANT columns use duckdb_variant.stream_canonical_arrow instead of the
plain Arrow reader below; it preserves the shredded representation and marker.
"""
from __future__ import annotations
from pathlib import Path
from uuid import uuid4
from raincloud import duckdb_connect
from ..canonical import open_canonical_writer
from ..spec import workdir_root
def my_streaming_handler(spec: dict, parsed: list[tuple[Path, object | None]]):
paths = [str(path) for path, _ in parsed if path.is_file()]
if not paths:
raise ValueError("streaming handler needs at least one input file")
# Uses the configured scratch disk and recipe-scoped directory. Unique
# scratch files let this handler clean up only resources it owns.
workdir = workdir_root() / spec["slug"]
workdir.mkdir(parents=True, exist_ok=True)
db_path = workdir / f"ingest-{uuid4().hex}.duckdb"
con = duckdb_connect(db_path)
try:
con.execute(
"CREATE TABLE records AS SELECT * FROM read_json_auto(?, format='newline_delimited')",
[paths],
)
reader = con.execute("SELECT * FROM records").to_arrow_reader(65536)
with reader, open_canonical_writer(spec["slug"], reader.schema) as writer:
for batch in reader:
writer.write_batch(batch)
finally:
con.close()
db_path.unlink(missing_ok=True)
Path(str(db_path) + ".wal").unlink(missing_ok=True)
return []