BigQueryDest
class BigQueryDest(
tables: BigQueryTableNamesSpec,
if_table_exists: IfTableExistsSpec = 'append',
schema_evolution: SchemaEvolutionSpec = 'strict',
merge_cfg: MergeCfgSpec | None = None,
watermark_column: BigQueryColumnSpec | None = None,
dest_cfg: DestCfgSpec | None = None,
)
Bases: Dest
Categories: destination
BigQuery destination for @subscriber -- one output slot per table.
Examples
Overwrite the target table each run with if_table_exists="replace":
@subscriber(
destination=BigQueryDest(
tables=["`cap_replace`"],
if_table_exists="replace",
),
input_tables=["caps_input/cap_rows"],
)
def publish(rows: TableFrameSpec) -> TableFrameSpec:
return rows
Append a wider frame; schema_evolution="update" adds the new column:
@subscriber(
destination=BigQueryDest(
tables=["`cap_evolve`"],
if_table_exists="append",
schema_evolution="update",
),
input_tables=["caps_input/cap_evo_v2"],
)
def evolve(v2: TableFrameSpec) -> TableFrameSpec:
return v2
Merge mode merges incoming rows onto the target by the keys in merge_cfg:
@subscriber(
destination=BigQueryDest(
tables=["`cap_merge`"],
if_table_exists="merge",
merge_cfg=[
MergeCfg(
on=["cap_id"],
op_column="op",
)
],
),
input_tables=["caps_input/cap_merge_v2"],
)
def merge(rows: TableFrameSpec) -> TableFrameSpec:
return rows
Parameters
tablesBigQueryTableNamesSpec (list[str])Target table names, one write slot per entry. Each is
[[project.]dataset.]table; a bare table or
dataset.table is completed from the connection's
project / dataset defaults. BigQuery wraps the whole
name in one backtick pair, e.g. my-project.ds.t.
if_table_existsIfTableExistsSpec (Literal['append', 'replace', 'merge'])"append" (the default) adds rows;
"replace" drops and recreates the table;
"merge" merges rows.
schema_evolutionSchemaEvolutionSpec (Literal['update', 'strict'])"strict" (the default) rejects schema changes;
"update" allows new columns.
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_columnBigQueryColumnSpec | 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 load job stamps the rows it
loads. A replace is refused it, having nothing to make idempotent --
it rewrites the table wholesale, so a rerun already leaves the same
rows behind.
dest_cfgDestCfgSpec | None (Mapping[Literal['bigquery.logging_level', 'bigquery.schema_evolution_engine'], Any] | None)Optional connector config. Supported keys:
bigquery.logging_level, and
bigquery.schema_evolution_engine -- "native" (the default
when unset) or "iceberg". Native evolution is BigQuery's
schema_update_options on the load job, which only ever adds a
column or relaxes one to nullable. "iceberg" has the connector
diff the target against the incoming batch instead, which also
widens a column whose type the batch has outgrown and refuses a
change BigQuery 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
bigquery.schema_evolution_engine="iceberg", since a MERGE takes no
load-job options 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 BigQueryValidateException.