Destination & load patterns¶
Prerequisites · A source block from Source patterns
and a dataflow envelope from Dataflows. Transform details are
optional for a simple load, but destination strategy can constrain them.
End state · A correct destination block with the right load_type and any
required supporting fields.
The Destination block decides where data is written and how existing data is treated. In DataCoolie, those are separate decisions:
- destination connection/format = physical target type
load_type= write semantics against that targetpartition_columnsandwrite_options= storage behavior modifiers
Destination at a glance¶
"destination": {
"connection_name": "silver",
"schema_name": "sales",
"table": "orders",
"load_type": "merge_upsert",
"merge_keys": ["order_id"],
"partition_columns": [
{ "column": "order_date", "expression": "CAST(created_at AS DATE)" }
],
"configure": {
"write_options": { "mergeSchema": true }
}
}
| Field | Required | Meaning |
|---|---|---|
connection_name |
conditional | Must match a destination connection; alternatively use inline connection |
schema_name |
no | Output namespace / folder |
table |
yes | Output table or folder name |
load_type |
no | append, overwrite, full_load, merge_upsert, merge_overwrite, scd2; defaults to append |
merge_keys |
conditional | Required for key-based merge and SCD2; a usable watermark replacement window can remove this requirement for merge_overwrite |
partition_columns |
no | Partition output by existing or derived columns |
configure.write_options |
no | Per-dataflow write-option overrides |
configure.merge_options |
no | Per-dataflow MERGE/upsert/SCD2 options; separate from append/overwrite writer options |
configure.scd2_effective_column |
conditional | Required for scd2 |
connection.configure.write_options and destination.configure.write_options
are merged, with destination overrides winning.
The same precedence applies to merge_options. Keep options that control the
MERGE condition (for example source_alias, target_alias, or predicate) in
merge_options; keep file/table writer options such as schema evolution or
compression in write_options. A legacy merge alias placed in
write_options is still read for compatibility, but it is not forwarded to
the append phase.
Which destination types are built in?¶
Built-in writers today support these destination families:
| Destination family | Supported formats | Addressing style | Supported load types | Maintenance |
|---|---|---|---|---|
| Flat-file output | parquet, csv, json, jsonl, avro |
Path-based | append, overwrite, full_load |
No |
| Lakehouse table | delta, iceberg |
Path-based or catalog/database/table | All registered load types, subject to engine/catalog support | Yes |
What is not built in
- Excel is not a writable destination
- database, API, and function destinations require custom destination plugins
Choose your strategy first¶
flowchart TD
A["Do you need full change history?"] -->|yes| SCD2["scd2"]
A -->|no| B["Can the source deliver the full current state?"]
B -->|"Full table every run"| OW["overwrite"]
B -->|"Only new/changed rows"| C["Are duplicates acceptable?"]
C -->|yes| APP["append"]
C -->|"No — need latest-wins upsert"| MU["merge_upsert"]
C -->|"No — need full refresh of matched keys"| MO["merge_overwrite"]
append — Add rows only, never touch existing¶
When to use: event streams, log ingestion, any scenario where duplicates in the destination are acceptable or impossible.
"destination": {
"connection_name": "bronze",
"schema_name": "events",
"table": "clicks",
"load_type": "append"
}
No extra fields required. Every run adds the rows returned by the source.
Works for both flat-file outputs and lakehouse tables.
For flat-file destinations, each append is written as a new file whose name contains the UTC write time and a short uniqueness suffix. This keeps rapid replay chunks and concurrent append calls from replacing one another; readers should scan the destination folder rather than depend on one fixed filename.
overwrite — Replace everything¶
When to use: daily snapshots, reference tables, aggregates that are always rebuilt from scratch.
"destination": {
"connection_name": "silver",
"schema_name": "sales",
"table": "daily_totals",
"load_type": "overwrite"
}
full_load is an alias
"load_type": "full_load" is equivalent. Use overwrite for new
metadata.
This is the safest whole-table strategy for file outputs. If your destination
format is parquet, csv, json, jsonl, or avro, overwrite / full_load
and append are the built-in options.
merge_upsert — Upsert by key (SCD1 / CDC)¶
When to use: incremental CDC-style loads, dimension tables that change over time but do not need history, customer/product master tables.
Rows matching merge_keys are updated; rows with no match are inserted.
Nothing is deleted.
"destination": {
"connection_name": "silver",
"schema_name": "sales",
"table": "customers",
"load_type": "merge_upsert",
"merge_keys": ["customer_id"]
}
| Field | Required | Notes |
|---|---|---|
merge_keys |
yes | List of column names that uniquely identify a row. Can be composite: ["order_id", "line_item_id"] |
Deduplication is often paired with merge_upsert
If your source can deliver duplicate rows for the same key, add
transform.deduplicate_columns and transform.latest_data_columns to
keep only the latest one before the merge. See
Transform patterns.
!!! info "First load falls back to overwrite" If the target table does not exist yet, DataCoolie performs an initial overwrite-style write and only uses merge semantics on later runs.
merge_overwrite — Rolling overwrite by key¶
When to use: nightly snapshot that always holds the full current state for a
rolling window. Existing target rows matching merge_keys are deleted then
re-inserted from the source. A disappeared source row with no incoming key is
not deleted by this key-based path; use a complete bounded replacement window
when upstream deletions must be reflected.
"destination": {
"connection_name": "silver",
"schema_name": "logistics",
"table": "active_shipments",
"load_type": "merge_overwrite",
"merge_keys": ["shipment_id"]
}
First load falls back to overwrite
Like merge_upsert, this strategy writes a brand-new table with overwrite
semantics when the destination does not exist yet.
replace_by_watermark — Range-based delete¶
When destination.configure.replace_by_watermark is true (see
destination.configure), the
merge_overwrite strategy switches from key-based delete to
range-based delete: instead of deleting only rows matching merge_keys,
it deletes all target rows within the watermark window and re-inserts
the full source batch. This handles upstream deletions that a key-match
would miss.
When the Driver has no usable replacement window, the strategy follows the
key-based path and requires merge_keys; replace_by_watermark does not invent
a window by itself. The source still needs a watermark column and an authored
look-back option for the range-based case.
Dataflow fragment for a complete source and destination combination:
{
"source": {
"connection_name": "shipments_source",
"table": "active_shipments",
"watermark_columns": ["updated_at"],
"configure": { "backward_days": 3 }
},
"destination": {
"connection_name": "silver",
"schema_name": "logistics",
"table": "active_shipments",
"load_type": "merge_overwrite",
"configure": { "replace_by_watermark": true }
}
}
Requirements:
- The source or its referenced connection must have an authored look-back
option such as
backward_daysorbackward. DataCoolie computes the effectivedate_backwardvalue at runtime; do not adddate_backwardto metadata. - Only supported with
merge_overwriteload type. - The source must cover the complete window, including rows deleted upstream.
merge_keysis needed when the strategy falls back to key-based overwrite; it is not required for a usable range replacement window.
At runtime the pipeline builds an immutable, attempt-local window from the
source watermark observations (or explicit replay bounds). The
MergeOverwriteStrategy passes that window to the engine's
replace_window(...) operation. The engine validates the final output
columns before deleting the bounded target rows, then appends the fresh batch;
the authored dataflow metadata is never mutated. A normal empty incremental
read skips the replacement, while a confirmed empty replay can replace its
explicit scope without creating an unrelated destination.
See the worked case, including first-run behavior and the distinction between normal empty reads and explicit empty replay, in Replace a complete watermark window. For the source-side look-back forms and connection/source precedence, see Incremental windows and look-back.
scd2 — Slowly Changing Dimension Type 2¶
When to use: dimension tables where you need full change history — e.g.
customers, products, employees where you need to know what value was
current at any past point in time.
SCD2 stores one row per version of each entity. The framework automatically adds three audit columns to every version row:
| Column | Meaning |
|---|---|
__valid_from |
When this version became current (copied from your date column) |
__valid_to |
When this version ended — NULL means it is still current |
__is_current |
true for the active version |
"destination": {
"connection_name": "gold",
"schema_name": "dims",
"table": "customer",
"load_type": "scd2",
"merge_keys": ["customer_id"],
"configure": { "scd2_effective_column": "updated_at" }
}
| Field | Required | Notes |
|---|---|---|
merge_keys |
yes | The natural/business key of the entity |
configure.scd2_effective_column |
yes | The source column that timestamps when this version became effective |
SCD2 tables grow over time
Each run appends new versions for changed rows. Plan your storage and run maintenance (vacuum/optimize) regularly.
Only send strictly newer versions
The close step ignores equal/older scd2_effective_column values, but the
append step inserts every incoming row. Filter and deduplicate upstream so
an older value cannot create a second open current version.
!!! info "First load falls back to overwrite" On the first run, DataCoolie creates the destination table first and then switches to SCD2 versioning on subsequent runs.
When reading a stored SCD2 table, select the current version or a historical instant explicitly. These are SQL query examples, not additional SCD2 metadata fields:
SELECT * FROM customers_scd2
WHERE __valid_from <= CAST('2026-01-01' AS TIMESTAMP)
AND COALESCE(__valid_to, CAST('9999-12-31' AS TIMESTAMP)) > CAST('2026-01-01' AS TIMESTAMP)
For an incremental source that feeds this strategy, see SCD2 with incremental inputs.
partition_columns — Partition the output table¶
Any destination strategy can optionally write partitioned data. This is not a
load type on its own — it is an addition to the destination block.
Each partition_columns item follows the Partition column
shape.
For example:
"destination": {
"connection_name": "silver",
"schema_name": "sales",
"table": "orders",
"load_type": "overwrite",
"partition_columns": [
{ "column": "order_date", "expression": "CAST(created_at AS DATE)" }
]
}
expression is evaluated as SQL before the write. The computed column
(order_date) is added to the DataFrame by PartitionHandler (transformer
order 80) and used as the partition key.
Use the top-level destination.partition_columns shape for new metadata. The
model also lifts destination.configure.partition_columns when the top-level
field is omitted. Database-backed metadata keeps a legacy
destination_configure.partition_by alias for existing rows and gives that
alias precedence over partition_columns; do not author both forms in a new
dataflow.
Partition by an existing column — no expression needed:
Multi-level partitioning:
"partition_columns": [
{ "column": "order_year", "expression": "EXTRACT(YEAR FROM order_date)" },
{ "column": "order_month", "expression": "EXTRACT(MONTH FROM order_date)" }
]
Partition columns extend merge keys internally
For merge-style destinations, DataCoolie automatically appends destination partition columns to the internal merge-key set when they are not already present. This keeps merge semantics aligned with the physical partitioning.
Use CAST, not date(...), for Polars portability
For partition expressions, prefer CAST(created_at AS DATE) over
date(created_at) so the metadata works in Polars as well as Spark.
Partition expression portability¶
PartitionHandler runs at order 80, after SystemColumnAdder at order 70,
so expressions can use __updated_at, for example
{"column": "etl_date", "expression": "CAST(__updated_at AS DATE)"}.
Spark evaluates expressions with Spark SQL; Polars uses pl.sql_expr, which
supports a smaller SQL subset. Prefer CAST(updated_at AS DATE) over a
provider-specific date(...) function. For year extraction use
EXTRACT(YEAR FROM updated_at). Avoid assuming current_timestamp() or
Java-style date_format(..., "yyyy-MM-dd") works in Polars; use the system
timestamp or portable CAST instead. Check every expression on the selected
engine before running it in another engine. Also see the
PartitionHandler execution order.
For merge-style writes, partition columns extend the effective merge key. If an existing entity changes partition value, a key-based merge may not match its old row. Choose a stable partition identity or explicitly account for relocation; see SCD2 with incremental inputs.
Flat-file outputs with date folders¶
Flat-file destinations have one more path-shaping option on the connection:
{
"name": "curated_parquet",
"connection_type": "file",
"format": "parquet",
"configure": {
"base_path": "data/output/curated",
"date_folder_partitions": "{year}/{month}/{day}"
}
}
This writes under a dated subfolder such as
data/output/curated/sales/orders/2026/05/09.
Use this when the partitioning is based on load time rather than a column in the DataFrame.
partition_columns wins over date_folder_partitions
For flat-file destinations, if you configure both, DataCoolie uses
partition_columns and ignores the date-folder pattern.
For overwrite and full_load, only the folder resolved for the current UTC
load time is replaced. Older dated folders remain available as immutable
snapshots; the operation does not delete the base path or sibling date folders.
append writes a new file in the resolved current folder.
Write options¶
Put write-engine options at the connection level when most dataflows should use
them, and in destination.configure
when only one dataflow needs them.
Connection-level defaults:
{
"name": "silver",
"connection_type": "lakehouse",
"format": "delta",
"configure": {
"base_path": "data/output/silver",
"write_options": {
"mergeSchema": true
}
}
}
Destination-level override:
"destination": {
"connection_name": "silver",
"schema_name": "sales",
"table": "orders",
"load_type": "append",
"configure": {
"write_options": {
"compression": "zstd"
}
}
}
Operation-specific options¶
Common writer options use the Spark-style vocabulary where an equivalent
meaning is qualified for the selected engine. Native options remain available
through the same maps. Use merge_options for the merge phase of
merge_upsert, merge_overwrite, and scd2:
"destination": {
"load_type": "merge_overwrite",
"merge_keys": ["order_id"],
"configure": {
"merge_options": {
"source_alias": "incoming",
"target_alias": "current"
},
"write_options": {
"schema_mode": "merge"
}
}
}
merge_options are sent only to the MERGE builder. write_options are sent to
the append/overwrite writer, including the append phase after a merge-overwrite
or SCD2 operation. Omitted options keep the existing engine defaults; this
section does not redefine those defaults. The canonical option inventory is
still being qualified per format, operation, and engine.
Named Polars Iceberg MERGE and SCD2 operations currently have no generic
options API. Supplying non-empty merge_options or write_options for those
operations fails before mutation until a backend mapping is qualified; this
prevents an option from being silently ignored.
Advanced lakehouse registration options¶
These options matter mostly for Delta/Iceberg deployments with metastore or AWS catalog integration. Their connection-level shape is under Connection settings by endpoint type:
| Option | Where | What it does |
|---|---|---|
catalog / database |
connection field or configure |
Registers or addresses the destination by qualified name instead of path only |
athena_output_location |
connection.configure |
After Delta writes and maintenance, registers a native Delta table through Athena DDL |
generate_manifest |
connection.configure |
Generates _symlink_format_manifest/ after writes and maintenance |
register_symlink_table |
connection.configure |
Registers a Glue symlink table; implies manifest generation |
symlink_database_prefix |
connection.configure |
Prefix for the generated symlink database name |
When catalog or database is present, DataCoolie identifies the physical
destination by qualified table name. Otherwise it identifies it by path.
That distinction matters for maintenance deduplication and fan-in orchestration.
AWS Delta registration: complete prerequisites¶
Use the AWS registration options only when the destination is a Delta path on a platform with the AWS integration enabled. A complete configuration includes:
{
"name": "orders_athena",
"connection_type": "lakehouse",
"format": "delta",
"configure": {
"base_path": "s3://analytics-lake/silver",
"catalog": "AwsDataCatalog",
"database": "analytics_silver",
"athena_output_location": "s3://analytics-query-results/datacoolie/",
"generate_manifest": true,
"register_symlink_table": true
}
}
The runtime must be using the AWS platform implementation, the Delta path must resolve, and the Athena output location and target database must be available to the configured credentials. Without those prerequisites registration is skipped or cannot complete; a successful local file write does not prove that the Glue/Athena catalog operation succeeded. Run a bounded smoke test and verify the registered table separately.
Maintenance support¶
- Flat-file destinations do not support maintenance.
- Delta and Iceberg destinations do.
- Maintenance is dispatched per physical destination, not per metadata row, so duplicate dataflows targeting the same table/path are deduplicated.
See User guide · Maintenance (vacuum/optimize) for the operational workflow.
Common mistakes¶
| Symptom | Likely cause | Fix |
|---|---|---|
FileWriter only supports ['append', 'full_load', 'overwrite'] |
Tried merge_upsert, merge_overwrite, or scd2 on a flat-file destination |
Use a Delta/Iceberg destination for merge semantics, or switch to append / overwrite |
merge_keys required error |
Used merge_upsert or scd2 without merge_keys |
Add "merge_keys": ["your_key_column"] |
| SCD2 columns not added | scd2_effective_column missing from configure |
Add "configure": { "scd2_effective_column": "updated_at" } in destination |
| Full table replaced when you wanted upsert | load_type is overwrite instead of merge_upsert |
Change load_type |
| Duplicate rows in destination after merge | Source delivers multiple rows for same key; dedup not configured | Add transform.deduplicate_columns — see Transform patterns |
| Partition expression fails on Polars | Used date(col) or another unsupported SQL function |
Use CAST(col AS DATE) or EXTRACT(...) |
| Maintenance skipped or fails on file outputs | Flat-file destinations do not implement maintenance | Run maintenance only on Delta/Iceberg destinations |