Add Dynamic Schema support to StorageWriteToBigQuery - #39236
Add Dynamic Schema support to StorageWriteToBigQuery#39236jrmccluskey wants to merge 15 commits into
Conversation
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request adds support for dynamic schemas and schema side inputs when writing to BigQuery using the Storage Write API (StorageWriteToBigQuery). It includes updates to handle callable schemas, passes schema side inputs through the pipeline, adds validation to ensure dynamic schemas are only used with dynamic destinations, and provides comprehensive unit and integration tests. There are no review comments, and the implementation looks solid, so I have no additional feedback to provide.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
There was a problem hiding this comment.
Code Review
This pull request introduces support for dynamic schemas when writing to dynamic destinations using the BigQuery Storage Write API, including a new dynamic_schema helper function, updates to StorageWriteToBigQuery to handle dynamic schemas and side inputs, and comprehensive unit tests. The reviewer feedback highlights critical issues that need to be addressed: potential runtime NameError exceptions from referencing the bigquery module namespace within bigquery.py, unintended in-place mutation of user-provided schemas during merging, and potential serialization/encoding errors (CoderException) alongside performance overhead caused by dynamic beam.Row instantiation with None values.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>
|
Caution The consumer version of Gemini Code Assist on GitHub has been sunset. All code review activity has officially ceased. |
|
Assigning reviewers: R: @shunping for label python. Note: If you would like to opt out of this review, comment Available commands:
The PR bot will only process comments in the main thread (not review comments). |
|
Reminder, please take a look at this pr: @shunping |
|
r: @claudevdm Could help to review this since you are more familiar with BigQueryIO? Thanks! |
|
Stopping reviewer notifications for this pull request: review requested by someone other than the bot, ceding control. If you'd like to restart, comment |
claudevdm
left a comment
There was a problem hiding this comment.
If I understand this PR correctly, the main restriction is that every xlang PCollection has exactly one element type/ beam schema that is encoded with RowCoder. Both SDK's need to agree on wire format.
If N destinations with N different schemas flow through one PCollection then the schema must be a superset of all of them.
Roughly:
| # | Stage | Element type | Code |
|---|---|---|---|
| 1 | user input | dict |
- |
| 2 | AppendDestinationsFn appends the destination |
(destination: str, dict) |
bigquery.py:2869, bigquery_tools.py:1581 |
| 3 | ConvertToBeamRows builds the Row |
beam.Row(destination=str, record=beam.Row(...)) |
bigquery.py:2961, applied at 2889 |
| 4 | .with_output_types(...) declares the PCollection schema → picks the RowCoder |
(type declaration) | bigquery.py:3049 |
| 5 | SchemaAwareExternalTransform crosses the boundary |
RowCoder-encoded bytes |
bigquery.py:2909 |
| 6 | expand(PCollectionRowTuple) |
PCollection<Row> |
BigQueryStorageWriteApiSchemaTransformProvider.java:177 |
| 7 | getFilterFormatFunction unwraps record |
TableRow |
PortableBigQueryDestinations.java:107 |
| 8 | TableRowConverter.toMessage |
protobuf message | StorageApiDynamicDestinationsTableRow.java:83 |
| 9 | AppendRows |
- | BigQuery |
Lets say there are two destinations with different schemas
p:d.users -> id:INTEGER, name:STRING
p:d.scores -> id:INTEGER, score:INTEGER, active:BOOLEAN
Stage 2 destination appended:
('p:d.users', {'id': 1, 'name': 'alice'})
Stage 3/4 the declared PCollection type is one schema, so it must cover both destinations:
Row(destination=<class 'str'>,
record=Row(id=Optional[numpy.int64],
name=Optional[str],
score=Optional[numpy.int64],
active=Optional[bool]))
Row(destination='p:d.users',
record=Row(id=1, name='alice', score=None, active=None))
Stage 7 getFilterFormatFunction(fetchNestedRecord=true) unwraps record and converts:
TableRow{id=1, name="alice", score=null, active=null}
BigQueryUtils.toTableRow emits every field including nulls:
output = output.set(schemaField.getName(), fromBeamField(schemaField.getType(), value));
Java
The only schema Java knows about is the schema of the PCollection (full union).
This means that every destination table actually has to comply with the full union schema? If they were pre-created with a subset of schemas writes should fail? If they are auto created in java, all tables will have the full union schema? What happens when merging different modes (required/nullable, repeated)? Some of this can potentially be fixed by making
The complexity and potential issues come from requiring a union schema.
Can we consider not forcing every row through one PCollection? How about we
- Group destinations by distinct schema in python
- Route each element to a branch matching its destination schema. Each branch has its own write path with a single concrete schema
- Dynamic destinations still apply within each branch (distinct table schema). Number of branches == number of disticnt schemas. 1000 tables sharing 2 schemas is 2 branches
- Since each PColl has one homogenous record schema, we avoid union, padding, merging and all that complexity.
# one branch per DISTINCT SCHEMA (not per table)
routed = elements | beam.ParDo(RouteBySchema()).with_outputs(*tags)
for schema, tag in tag_of_schema.items():
routed[tag] | beam.io.WriteToBigQuery(
table=get_destination, # still dynamic destinations within the branch
schema=schema, # ...but ONE concrete schema
method=WriteToBigQuery.Method.STORAGE_WRITE_API)
Inside a branch all destinations genuinely share the schema, so getSchema() returning one schema is correct, no Java change needed.
| Dynamic Schemas with Storage Write API | ||
| -------------------------------------- | ||
| When writing to dynamic destinations with `method=STORAGE_WRITE_API`, a union schema | ||
| containing all fields across destination tables is required at the PCollection level |
There was a problem hiding this comment.
Is the expectation that the destination schema for every table will be the full union schema, or each table is meant to have the individually mapped schema and union is just a trick to make it work? Can you add an integration test that validates the expected schema on the tables and document the expected behavior?
| use_at_least_once=False)) | ||
| hamcrest_assert(p, all_of(*bq_matchers)) | ||
|
|
||
| def test_write_to_dynamic_destinations_with_dynamic_schema(self): |
There was a problem hiding this comment.
Can we add a test for the case where the destination tables are pre-created/ already exist beforehand?
| { | ||
| "comment": "Modify this file in a trivial way to cause this test suite to run", | ||
| "modification": 16 | ||
| "modification": 21 |
There was a problem hiding this comment.
I think you can also trigger .github/trigger_files/beam_PostCommit_Python_Xlang_Gcp_Direct.json for faster validation
| schema_side_inputs=(beam.pvalue.AsSingleton(schema_pc), ), | ||
| use_at_least_once=False)) | ||
| hamcrest_assert(p, all_of(*bq_matchers)) | ||
|
|
There was a problem hiding this comment.
Can we also validate the expected schema of the tables that were created after the pipeline ran?
Implements support for using dynamic schemas in StorageWriteToBigQuery invocations.
Thank you for your contribution! Follow this checklist to help us incorporate your contribution quickly and easily:
addresses #123), if applicable. This will automatically add a link to the pull request in the issue. If you would like the issue to automatically close on merging the pull request, commentfixes #<ISSUE NUMBER>instead.CHANGES.mdwith noteworthy changes.See the Contributor Guide for more tips on how to make review process smoother.
To check the build health, please visit https://github.com/apache/beam/blob/master/.test-infra/BUILD_STATUS.md
GitHub Actions Tests Status (on master branch)
See CI.md for more information about GitHub Actions CI or the workflows README to see a list of phrases to trigger workflows.