Operations examples¶
Operational samples use the same Driver lifecycle as a normal run. Keep one
job identity for a Driver session, pass explicit external run_attributes, and
let the Driver decide result status and teardown ownership.
Incremental project recipe¶
The Incremental project is the smallest local
stateful recipe. It requires Python 3.11+ and
datacoolie[polars]==0.2.0; the installation guide
documents the matching source-wheel handoff when that release is not
published. For that preview/source-wheel handoff, apply the polars extra to
the matching wheel for a minimal install; the guide's generic
cli,polars-delta profile already supplies Polars for this recipe. Extract
incremental.zip,
change to the extracted incremental/ root, and run the project-owned runner:
python -m pip install "datacoolie[polars]==0.2.0"
python runners/dev/run.py --state-base-path .runtime
The archive starts with two CSV rows (updated_sequence 1 and 2). The first
run writes two business rows under data/output/orders/ and saves the source
watermark below .runtime/watermarks/. Run the same command again without
changing data/input/orders/orders.csv; the persisted watermark makes this a
successful no-change run and the output remains two business rows.
Append one later source row, then run the same command a third time:
from pathlib import Path
source = Path("data/input/orders/orders.csv")
with source.open("a", encoding="utf-8") as handle:
handle.write("3,5.50,3\n")
Read the append output and watermark independently:
import json
from pathlib import Path
import polars as pl
files = sorted(Path("data/output/orders").glob("*.parquet"))
rows = pl.concat([pl.read_parquet(path) for path in files]).select(
"order_id", "amount", "updated_sequence"
).sort("order_id")
assert rows.rows() == [(1, 19.99, 1), (2, 29.0, 2), (3, 5.5, 3)]
watermark = next(Path(".runtime/watermarks").rglob("watermark_value.json"))
assert json.loads(watermark.read_text(encoding="utf-8"))[
"updated_sequence"
] == 3
This project uses an append destination, so rerunning after a failed write can
repeat committed business rows. To adapt it, append source rows with a value
greater than the stored watermark and update watermark_columns and schema
hints together when the key changes. To reset a trial, extract a fresh copy of
the ZIP into a new directory; generated output and .runtime/ belong to each
extracted copy and are not removed by the recipe.
Replay¶
runners/local/replay.py (source ·
raw) demonstrates a bounded
[start, end) interval, optional chunks and an explicit confirmation before a
watermark is saved. The managed-host equivalent is
runners/databricks/replay_spark.ipynb
(source ·
raw).
The Incremental project (project-files ·
download) provides the input snapshot used by this
replay recipe. The ZIP contains the project runner but not the operation
wrapper. Obtain the separate raw runners/local/replay.py file
(source ·
raw) and save it inside the extracted
incremental/ project root as replay.py.
Start this replay lesson from a fresh extraction of incremental.zip. The
incremental recipe above deliberately leaves row 3, output files and a saved
watermark in place; reusing that workspace changes the replay input and result
counts.
python replay.py --metadata-path metadata --watermark-base-path .runtime/watermarks --log-base-path .runtime/logs --working-directory . --start 1 --end 4
Run that command from the extracted incremental/ root. The runner changes to
--working-directory before constructing FileProvider, so metadata, the
watermark root and the log root all resolve inside that project. Relative paths
from a repository checkout such as docs/examples/files/... do not remain
valid after that directory change. On the fresh two-row snapshot this writes
two business rows; because --save-watermark is absent, it does not persist a
replay watermark.
Restartable replay and interrupted sessions¶
operations/replay_recovery.py (source ·
raw) is a small
project-owned wrapper around the same public replay API. Run the requested
range again with the same input snapshot and watermark root after a failed or
interrupted session. Obtain this raw wrapper separately and save it as
replay_recovery.py in the extracted project root:
Use a fresh extracted project for the first recovery attempt, then repeat the same command in that workspace when checking retry behavior. Every repeat runs the requested chunks again.
python replay_recovery.py --metadata-path metadata --watermark-base-path .runtime/watermarks --log-base-path .runtime/logs --working-directory . --start 1 --end 4 --chunk-interval-json '{"step": 2}' --save-watermark --confirm-save-watermark --job-id replay-attempt-2
With the two-row snapshot and step=2, the first recovery run writes the two
requested business rows and saves updated_sequence=2. Every requested chunk
runs again on the next invocation; an append destination can therefore contain
two copies of each row after the repeat. save_watermark persists the reader's
observations after successful writes but is not a replay checkpoint. The
destination write and watermark write are separate operations, so a hard
termination after the destination commit can still leave output to be written
again. The framework does not promise exactly-once delivery for an append
target; use a supported keyed strategy such as merge_upsert when retries must
reconcile rows. Keep each attempt's job_id distinct and retain its logs for
diagnosis.
The runner does not inject failures. Failure-boundary checks belong in an isolated test process so a public project script never carries a production fault-injection switch.
Maintenance¶
runners/local/maintenance.py (source ·
raw) requires an
explicit --confirm-maintenance flag and refuses a no-op request. Maintenance
deduplicates by physical destination before dispatching compact/cleanup work.
Sharding and failure behavior¶
--job-num and --job-index identify a disjoint shard of the selected
dataflows. The caller owns any external stage barrier: wait for every job before
starting a dependent stage. A failed result is returned and a runner should exit
non-zero; preparation failures are reported before business reads are retried.
Logging source files¶
The persisted structured records are full JSON objects. Snapshot mode rewrites a
stable job/dataflow projection; batch mode emits valid JSON Lines records in
.json files. The console can be colored for modern terminals, but color is
presentation-only.
Use configuration/logging_modes.py (source · raw) for the LogConfig sample, and operations/replay_recovery.py (source · raw) for its wrapper source. The normative contracts are logging and logging reference.
Do not call a logger's private flush method from a runner. Configure the public
LogConfig, close the Driver in a context manager or finally, and preserve the
structured records for diagnostics.