SnowflakeDest
class SnowflakeDest(
tables: SnowflakeTableNamesSpec,
if_table_exists: IfTableExistsSpec = 'append',
schema_evolution: SchemaEvolutionSpec = 'strict',
merge_cfg: MergeCfgSpec | None = None,
watermark_column: SnowflakeColumnSpec | None = None,
dest_cfg: DestCfgSpec | None = None,
)
Bases: Dest
Categories: destination
Snowflake destination for @subscriber -- one output slot per table.
Examples
Append rows to a table -- the default disposition, which accumulates across runs:
@subscriber(
destination=SnowflakeDest(
tables=["cap_append"],
if_table_exists="append",
),
input_tables=["caps_input/cap_rows"],
)
def load_append(rows: TableFrameSpec) -> TableFrameSpec:
return rows
Replace mode truncates the target before COPY INTO,
preserving the existing schema on reload:
@subscriber(
destination=SnowflakeDest(
tables=["cap_replace"],
if_table_exists="replace",
),
input_tables=["caps_input/cap_rows"],
)
def load_replace(rows: TableFrameSpec) -> TableFrameSpec:
return rows
Merge mode upserts by key instead of accumulating, and applies
the CDC tombstones the incoming op column marks:
@subscriber(
destination=SnowflakeDest(
tables=["cap_merge"],
if_table_exists="merge",
merge_cfg=[MergeCfg(on=["cap_id"], op_column="op")],
dest_cfg={"snowflake.schema_evolution_engine": "iceberg"},
),
input_tables=["caps_input/cap_changes"],
)
def load_merge(changes: TableFrameSpec) -> TableFrameSpec:
return changes
Parameters
tablesSnowflakeTableNamesSpec (list[str])Target table names, one write slot per entry. Each is a
1-, 2- or 3-part name (table, schema.table or
database.schema.table); double-quote a part for
case-sensitive / special-char names.
if_table_existsIfTableExistsSpec (Literal['append', 'replace', 'merge'])"append" (the default) adds rows;
"replace" issues TRUNCATE before COPY INTO (so
the existing schema is preserved);
"merge" stages the rows and MERGEs them onto the
target by the keys in merge_cfg, updating, deleting or
inserting each one.
schema_evolutionSchemaEvolutionSpec (Literal['update', 'strict'])"strict" (the default) rejects schema changes;
"update" adds new columns to the target to match the schema of
the incoming data.
merge_cfgMergeCfgSpec | None (list[MergeCfg] | None)Required by -- and only valid with --
if_table_exists="merge": how incoming rows are matched
onto existing ones. A list, one entry per table in tables and
in the same order -- a single table takes a one-entry list.
watermark_columnSnowflakeColumnSpec | None (str | None)An optional column stamped with the id of the trx the
write belongs to, making the write idempotent: a rerun that finds
the target already carrying this trx's id writes nothing rather than
applying the same batch twice. Every table in tables is stamped in
the same column. The target gains the column under either
schema_evolution setting -- a watermark no row can be stamped in
would do nothing at all -- and rows written before it existed keep a
NULL there.
A merge stamps through the MERGE itself; an append writes the
column into the batch it stages, so the COPY stamps the rows it
loads. A replace is refused it, having nothing to make idempotent --
it truncates and reloads, so a rerun already leaves the same rows
behind.
dest_cfgDestCfgSpec | None (Mapping[Literal['snowflake.logging_level', 'snowflake.schema_evolution_engine'], Any] | None)Optional connector config. Supported keys:
snowflake.logging_level, and
snowflake.schema_evolution_engine -- "native" (the default
when unset) or "iceberg". Native schema evolution is performed with
Snowflake's
ENABLE_SCHEMA_EVOLUTION, which COPY INTO honours and which
only ever adds a column. "iceberg" diffs the
target against the incoming file instead, which also widens a
column whose type the batch has outgrown and refuses a change
Snowflake could not make. It applies only where the schema is
allowed to move, so it does nothing under
schema_evolution="strict"; if_table_exists="merge" requires
snowflake.schema_evolution_engine="iceberg", since a MERGE does not read
the native mark at all.
Methods
validatedef validate()
Cross-field validation: merge_cfg and if_table_exists="merge"
require each other, a list of configs must line up with tables,
watermark_column is not a replace's to name, and merging requires
schema_evolution_engine="iceberg".
Called by the framework; raises SnowflakeValidateException.