Skip to main content
GuideServerConfigure and deploy Tabsdata servers on your machine.TutorialsConfigure data integration workflows within a running Tabsdata server.Advanced TutorialsBuild end-to-end workflows between two specific systems.API ReferenceCLI ReferenceRelease Notes
Version: 2.1.0

BigQueryDest

class
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​

parameter

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.

parameter
if_table_existsIfTableExistsSpec (Literal['append', 'replace', 'merge'])

"append" (the default) adds rows; "replace" drops and recreates the table; "merge" merges rows.

parameter
schema_evolutionSchemaEvolutionSpec (Literal['update', 'strict'])

"strict" (the default) rejects schema changes; "update" allows new columns.

parameter
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.

parameter
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.

parameter
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​

method
validate
def 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.