fix(bigquery): give added tables a snapshot and a CDC checkpoint - #4769
Open
dtunikov wants to merge 119 commits into
Open
fix(bigquery): give added tables a snapshot and a CDC checkpoint#4769dtunikov wants to merge 119 commits into
dtunikov wants to merge 119 commits into
Conversation
Capture T (BigQuery's own CURRENT_TIMESTAMP(), not local wall-clock) once per
ExportTxSnapshot call and append FOR SYSTEM_TIME AS OF TIMESTAMP('<T> UTC') to
every table's EXPORT DATA statement, so all tables in a snapshot read a
consistent point in time.
Persist T as the initial CDC checkpoint via SetLastOffset, but from
ExportTxSnapshot itself rather than SetupReplication as the plan originally
sketched: SnapshotFlowWorkflow only calls ExportTxSnapshot when
InitialSnapshotOnly is true, and only calls SetupReplication when it's false
(cloneTablesWithSlot, or the no-snapshot CDC-only branch) - the two are
mutually exclusive per run, so T is never in scope inside SetupReplication.
Gating the write on !InitialSnapshotOnly here is therefore a forward-looking
no-op until a later chunk changes how BigQuery mirrors that continue into CDC
get their initial load wired up; documented in code at the capture site.
Refactored the SQL-building half of bigQueryExportQueryStatement into a pure
buildBigQueryExportSQL taking an already-resolved schema, so it's unit
testable without a live BigQuery client (mirrors the existing
bigQuerySchemaToQRecordSchema/datasetTable test patterns in this package).
…t code ExportTxSnapshot only ever runs on SnapshotFlowWorkflow's pure snapshot-only branch (InitialSnapshotOnly && DoInitialSnapshot) - the continue-to-CDC and CDC-only branches call SetupReplication instead and never touch ExportTxSnapshot. So the checkpoint-persist code added there, gated on !InitialSnapshotOnly, could never actually run; removed. SetupReplication now does what MySQL's does: capture a starting position (BigQuery's current timestamp T, no replication slot to open) and persist it as the initial CDC checkpoint. Unlike MySQL, BigQuery's initial load reads pre-exported Parquet from GCS rather than querying the live table, so when the mirror wants an initial load (req.DoInitialSnapshot), this also runs that export as of T - the per-table export-job loop is factored out of ExportTxSnapshot into a shared exportTablesAsOf helper so both callers use it. Returns a zero-value SetupReplicationResult (no slot/snapshot name), same as MySQL: cloneTablesWithSlot falls back to the mirror's configured SnapshotStagingPath when no override is given, which is exactly where this export writes.
…bject_pull_test.go
Newline/tab padding made the query harder to read; collapse to single-line format string.
Introduces a connector-specific mirror-config extension point on FlowConnectionConfigs/FlowConnectionConfigsCore, mirroring how Peer already does this for peer-level config. BigqueryCdcConfig (with a cdc_mode of APPENDS or CHANGES) is the first variant. No converter work is needed: flow/proto_conversions copies between the two messages by field number, so a oneof - whose wrapper types are scoped per parent message - is carried over without special-casing. No behavior change yet - BigqueryCdcConfig is not read anywhere. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
start_peer_flow_job builds FlowConnectionConfigs as a fully exhaustive struct literal (no ..Default::default()), so it needs every field set explicitly. The new source_connector_config oneof broke this build; nexus doesn't set any connector-specific mirror config today, so None is correct here.
Copying by field number handles the oneof without special-casing, but the wrapper types are scoped per parent message, so assert the variant is rebuilt against the destination's own type in both directions - and that an unset oneof stays unset rather than spuriously populating a variant. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…ecting Snapshot-only mirrors skip these checks (mirrors MySQL's ValidateMirrorSource pattern). CHANGES mode needs a real PK constraint plus enable_change_history on the source table; APPENDS mode needs an explicit MergeTree destination engine when there's no PK, since ORDER BY tuple() on a keyless ReplacingMergeTree collapses the table on writes.
…GE_TREE CH_ENGINE_REPLICATED_REPLACING_MERGE_TREE is the same collapsing dedup engine as the plain variant, just wrapped for replication (see how normalize.go's engine switch groups the two under one case) - a keyless table hits the same ORDER BY tuple() collapse either way, so the APPENDS-mode keyless-engine check needs to catch both.
Gives BigQueryConnector real CDCPullConnectorCore bodies (SetupReplConn, UpdateReplStateLastOffset, PullFlowCleanup, EnsurePullability) and a real PullRecords: self-paced polling of APPENDS() per mapped table over (checkpoint, upper], converting rows to InsertRecords via a new BigQuery value -> QValue converter, advancing the checkpoint text once the window closes. No delete/insert pairing yet -- that's CHANGES mode, next chunk.
CHANGES() reports an UPDATE as a delete+insert pair sharing the same PK and _CHANGE_TIMESTAMP, so PullRecords needs to pair those back into one UpdateRecord instead of emitting two records -- otherwise every update on a CHANGES-mode mirror would be replicated as a delete followed by an unrelated insert. Reads cdc_mode once per PullRecords call (it's a per-mirror setting) to pick APPENDS or CHANGES per the "Decisions locked in" plan.
- cleanup comments
…CHANGES bq functions
…st does not exist" error
Rewrites the now-stale Test_BigQuery_Source_CDC_Not_Supported, which asserted CDC gets rejected - chunk 3 replaced that with mode-specific validation. Adds e2e coverage for snapshot->CDC handoff, APPENDS insert-only polling, CHANGES insert/update/delete pairing, and pause/resume from the persisted checkpoint; none of this has been run against a live BigQuery instance yet.
- cleanup comments
bigquery.go holds the shared BigQueryConnector struct used by both source and destination code, so source-only changes there were wrongly picked up by the deprecated-connector labeler.
remove cdc_batches updates from isolated cdc branch
…_replication_state should be present after RecordTableReplicationAttempt is called
A table addition runs as a child flow carrying only the tables being added, and it never runs SetupReplication -- an InitialSnapshotOnly flow skips it. On the query-based CDC path that left added tables broken twice over. BigQuery's ExportTxSnapshot read its table list from the catalog, which still holds the mirror's pre-addition config while the addition runs, so it exported the mirror's existing tables and never the added ones. The snapshotting flow's own table mappings now travel through MaintainTx into ExportTxSnapshot instead. Nothing then wrote the added tables' CDC checkpoint. Their state row was first created by RecordQueryCDCAttempt with an empty cursor, and PullTableRecords re-seeds start from its own clock whenever the cursor is empty, which puts the whole window inside the safety lag: the window is never valid, the cursor never advances, and the table never replicates. ExportTxSnapshot now seeds each exported table's checkpoint the same way SetupReplication does, through a shared initializeQueryCDCCheckpoints. Covered by Test_BigQuery_CDC_Table_Addition_Mid_CDC, which times out waiting for the added table's snapshot without the first fix and finds an empty cursor_text without the second. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
dtunikov
commented
Sep 2, 2026
| defer connClose(ctx) | ||
|
|
||
| exportSnapshotOutput, tx, err := conn.ExportTxSnapshot(ctx, flowName, env) | ||
| exportSnapshotOutput, tx, err := conn.ExportTxSnapshot(ctx, flowName, tableMappings, env) |
Contributor
Author
There was a problem hiding this comment.
ExportTxSnapshot should use the mappings that are passed via state, not those that are stored in the catalog. Catalog table mapping are updated only after MaintainTx succeed.
I'll need to check if it can break existing running temporal activities since it changes the set of input parameters.
…table-addition-checkpoint # Conflicts: # flow/activities/flowable.go # flow/activities/flowable_query_cdc.go # flow/connectors/bigquery/qrep_object_pull.go # flow/connectors/clickhouse/object_sync.go # flow/connectors/clickhouse/query_cdc_sync.go # flow/connectors/clickhouse/table_function.go # flow/connectors/external_metadata/query_cdc_replication_state.go # flow/connectors/utils/stream.go # flow/connectors/utils/typed_cdc_stream_test.go # flow/e2e/bigquery_cdc_test.go # flow/model/model.go
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What changed
Fix table additions for running query-based CDC mirrors. The child addition flow does not run
SetupReplication, and the catalog still contains the old table mappings while it runs.TableMappingstoExportTxSnapshotso it snapshots the added tables.snapshotTime + 1us, using the same helper asSetupReplication.Testing
Added
Test_BigQuery_CDC_Table_Addition_Mid_CDC, covering the added table's initial snapshot, checkpoint creation, continued CDC for both tables, and mirror status.Resolves: DBI-1091