Replay & backfill¶
Prerequisites · Familiarity with watermark replay behaviour and a running DataCoolieDriver.
End state · A completed historical replay of a bounded range, with an explicit choice about whether the source watermark is updated.
driver.run_replay() re-processes a bounded window of historical data in
sequential chunks. By default it leaves the production watermark untouched;
save_watermark=True persists the reader-produced candidate after each
successful destination write. The candidate can represent an observed maximum
or a source-confirmed request end. Every invocation still re-runs the
requested range; a saved watermark is not a replay checkpoint. Use it to:
- Backfill a new destination table from historical source data.
- Repair a corrupt range after a source system issue.
- Re-run a date range after a schema or logic change.
The ReplayConfig reference
lists the fields and defaults. If a runner injects its own watermark manager,
use the watermark API for provider-backed
read/write contracts and the Driver API
for run_replay and constructor arguments.
Quick start¶
from datacoolie.orchestration.driver import DataCoolieDriver
from datacoolie.core.models.run_config import ReplayConfig
replay = ReplayConfig(
start="2025-01-01", # inclusive lower bound
end="2025-04-01", # exclusive upper bound
chunk_interval={"months": 1},
)
with DataCoolieDriver(engine=engine, platform=platform, metadata_provider=metadata) as driver:
dataflows = driver.load_dataflows(stage="bronze2silver")
result = driver.run_replay(dataflows=dataflows, replay=replay)
print(f"Dataflows: succeeded={result.succeeded}, failed={result.failed}")
Replay applies the same activation check as a normal run, including when dataflows are supplied directly. An inactive dataflow or source/destination connection produces one skipped outer result; no chunks or watermark operations start for it.
This replays January, February, and March 2025 as three independent chunks:
[Jan 1, Feb 1), [Feb 1, Mar 1), [Mar 1, Apr 1).
ReplayConfig fields¶
| Field | Type | Default | Description |
|---|---|---|---|
start |
str, date, datetime, or int |
— | Inclusive lower bound of the replay range |
end |
str, date, datetime, or int |
— | Exclusive upper bound of the replay range |
chunk_interval |
Dict[str, int] or None |
None |
Chunking interval (see below). None = single-shot replay |
save_watermark |
bool |
False |
When True, saves the reader-produced watermark candidate after each successful chunk. It does not skip chunks on a later invocation |
chunk_column |
str or None |
None |
Column used to select chunks. It may be independent of source.watermark_columns when the source reader supports bounded reads |
Range convention: [start, end)¶
The range is left-closed, right-open:
startis included in the first chunk.endis excluded (the first instant NOT replayed).
This aligns with Python range(), Spark partitioning, and ISO 8601 interval conventions.
chunk_interval keys¶
| Key | Alignment | Example |
|---|---|---|
years |
Calendar years | {"years": 1} |
months |
Calendar months | {"months": 1} or {"months": 3} |
weeks |
ISO weeks (Mon–Sun) | {"weeks": 1} |
days |
Calendar days | {"days": 1} |
hours |
Hours | {"hours": 6} |
minutes |
Minutes | {"minutes": 30} |
step |
Integer step | {"step": 10000} |
Use step when watermark_columns holds a row-number or integer sequence id
instead of a timestamp.
# Integer watermark — replay rows 0 to 1,000,000 in batches of 100k
ReplayConfig(
start=0,
end=1_000_000,
chunk_interval={"step": 100_000},
)
chunk_interval=None is also valid for a finite numeric range: it performs one
bounded read using the exact numeric [start, end) values. step creates
sequential integer chunks. Calendar keys (years, months, weeks, days,
hours, and minutes) create temporal chunks and require a date/datetime
compatible source column; they are not interchangeable with integer stepping.
Single-shot replay
Set chunk_interval=None (the default) to replay the entire range as one
operation. Useful for small ranges or when chunking is not needed.
Chunk column resolution¶
DataCoolie auto-resolves the chunk column from dataflow.source.watermark_columns[0].
If your dataflow has multiple watermark columns and the first is not the
column you want to chunk on, set chunk_column explicitly. Database,
lakehouse, file, and function readers can use a source-supported bounded-read
column that is not persisted as a watermark. API readers can do the same when
the selected field has a canonical range_param_mapping entry with both a
lower and an upper binding. The API mapping is the endpoint contract for the
selected field; it does not add that field to source.watermark_columns.
For example, an API can select a historical created_at range while keeping
updated_at as the only persisted watermark:
{
"watermark_columns": ["updated_at"],
"configure": {
"range_param_mapping": {
"created_at": {
"lower": {"name": "created_from", "operator": ">="},
"upper": {"name": "created_to", "operator": "<"},
"format": "iso",
"response_column": "created_at",
"watermark_value": "observed_max"
}
}
}
}
save_watermark=True still persists only the source watermark columns. The API
response must include an updated_at observation for that key to advance; a
slice that returns only created_at rows has no updated_at value to save. A
saved updated_at maximum from a created_at slice is an observed value, not
proof that every updated_at value in the replay interval was covered.
ReplayConfig(
start="2025-01-01",
end="2025-04-01",
chunk_interval={"months": 1},
chunk_column="event_date", # override; uses this column instead of watermark_columns[0]
)
Persisting the reader watermark candidate (save_watermark=True)¶
By default (save_watermark=False), the production watermark is never touched
during replay. This is safe for backfill into a new table.
When save_watermark=True, DataCoolie saves the reader's candidate after each
successful destination write. For an observed-max binding, the candidate is the
maximum non-null source value; for a request-end binding, it is the exclusive
end that the source contract confirms. The requested chunk upper bound is never
substituted for a missing observation. An empty read or an all-null candidate
does not save state, while a legitimate zero remains a valid observation. The
driver merges that candidate with stored state using the reader's
source-qualified ordering semantics before it saves. On a later invocation the
same [start, end) range is read again:
ReplayConfig(
start="2025-01-01",
end="2025-07-01",
chunk_interval={"months": 1},
save_watermark=True, # persist source observations per successful chunk
)
Coordinate with normal incremental runs
When save_watermark=True, a successful chunk may update production state.
Source-qualified ordered keys never move backward: a historical candidate
below a higher saved value leaves that value unchanged. Coordinate replay
with normal incremental runs that share the same dataflow state.
Saving the maximum of a partially loaded history does not prove that every earlier record has been loaded. Complete the intended historical coverage before relying on that state for normal incremental selection.
No chunk checkpoint
Chunks run sequentially and later chunks stop after the first failure.
A restart intentionally runs every requested chunk again. Use an idempotent
destination strategy such as merge_upsert when repeated replay writes
must reconcile business rows.
Destination write and state save are separate
A successful destination write happens before the watermark save. A process termination or storage error in that gap can leave committed output with the old watermark. Rerun the complete requested range; use a keyed or otherwise idempotent destination strategy when duplicate delivery would be harmful. A later replay never treats the saved watermark as a completed chunk marker.
Combining with source.filter_expression¶
If your source dataflow already has a source.filter_expression, the source
reader keeps it as a separate post-read or pushdown filter. A replay range is
owned by the source reader and is applied independently; your original filter
is preserved:
"source": {
"connection_name": "orders_db",
"table": "orders",
"watermark_columns": ["updated_at"],
"filter_expression": "status = 'active'"
}
For a database source, the effective predicate for each chunk is:
The range is always [start, end). A date-folder file partition may be used
to prune folders, but its folder key is internal and cannot be a
chunk_column; file modification time remains the selector for files.
CLI boundary values¶
The downloadable replay_recovery.py
wrapper accepts --start and --end as text. A string made only of an optional
sign and decimal digits becomes an integer; every other value remains a string
for the source's date/datetime parser. The wrapper does not parse arbitrary
floating-point text or attach a timezone to a date string. --chunk-interval-json
must decode to an object. Saving state is an explicit opt-in: pass both
--save-watermark and --confirm-save-watermark; omitting them uses
save_watermark=False.
Execution model¶
- Dataflows are processed concurrently (bounded by
max_workers). - Chunks within a single dataflow run sequentially — a failed chunk stops further chunks for that dataflow but does not affect other dataflows.
- With
stop_on_error=True, no new replay dataflow is admitted after a terminal failure; already-admitted dataflows finish and withheld dataflows remain pending. - Each chunk records a separate
DataFlowRuntimeInfoin the execution log withoperation_type = "replay". - Rows written by a chunk receive that chunk's
__dataflow_run_id; they do not receive the outer replay aggregate ID. - Each chunk produces a separate execution log row.
- Chunk observers use the chunk completion hook; the dataflow completion hook receives one aggregate result. Observer errors are logged and do not change pipeline status.
ExecutionResult.total/succeeded/failedcount the outer dataflows, because_process_replayaggregates its chunks into oneDataFlowRuntimeInfobefore returning toParallelExecutor.
Full example: one-year backfill¶
from datetime import date
from datacoolie.orchestration.driver import DataCoolieDriver
from datacoolie.core.models.run_config import ReplayConfig
replay = ReplayConfig(
start=date(2024, 1, 1),
end=date(2025, 1, 1), # full year 2024
chunk_interval={"months": 1},
save_watermark=False, # backfill: leave production watermark untouched
)
with DataCoolieDriver(
engine=engine,
platform=platform,
metadata_provider=metadata,
log_base_path="logs/",
) as driver:
dataflows = driver.load_dataflows(stage="bronze2silver")
result = driver.run_replay(dataflows=dataflows, replay=replay)
print(f"Total dataflows: {result.total}")
print(f"Succeeded: {result.succeeded}")
print(f"Failed: {result.failed}")