Skip to content

Add Dynamic Schema support to StorageWriteToBigQuery - #39236

Open
jrmccluskey wants to merge 15 commits into
apache:masterfrom
jrmccluskey:dynamicDuo
Open

Add Dynamic Schema support to StorageWriteToBigQuery#39236
jrmccluskey wants to merge 15 commits into
apache:masterfrom
jrmccluskey:dynamicDuo

Conversation

@jrmccluskey

Copy link
Copy Markdown
Contributor

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:

  • Mention the appropriate issue in your description (for example: 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, comment fixes #<ISSUE NUMBER> instead.
  • Update CHANGES.md with noteworthy changes.
  • If this contribution is large, please file an Apache Individual Contributor License Agreement.

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)

Build python source distribution and wheels
Python tests
Java tests
Go tests

See CI.md for more information about GitHub Actions CI or the workflows README to see a list of phrases to trigger workflows.

@jrmccluskey

Copy link
Copy Markdown
Contributor Author

/gemini review

@apache apache deleted a comment from gemini-code-assist Bot Jul 8, 2026
@apache apache deleted a comment from gemini-code-assist Bot Jul 8, 2026

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread sdks/python/apache_beam/io/gcp/bigquery.py
Comment thread sdks/python/apache_beam/io/gcp/bigquery.py
Comment thread sdks/python/apache_beam/io/gcp/bigquery.py
Comment thread sdks/python/apache_beam/io/gcp/bigquery.py
jrmccluskey and others added 3 commits July 21, 2026 11:05
Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>
@jrmccluskey jrmccluskey changed the title [WIP] Add Dynamic Schema support to StorageWriteToBigQuery Add Dynamic Schema support to StorageWriteToBigQuery Jul 21, 2026
@jrmccluskey
jrmccluskey marked this pull request as ready for review July 21, 2026 20:05
@gemini-code-assist

Copy link
Copy Markdown
Contributor

Caution

The consumer version of Gemini Code Assist on GitHub has been sunset. All code review activity has officially ceased.

@github-actions

Copy link
Copy Markdown
Contributor

Assigning reviewers:

R: @shunping for label python.

Note: If you would like to opt out of this review, comment assign to next reviewer.

Available commands:

  • stop reviewer notifications - opt out of the automated review tooling
  • remind me after tests pass - tag the comment author after tests pass
  • waiting on author - shift the attention set back to the author (any comment or push by the author will return the attention set to the reviewers)

The PR bot will only process comments in the main thread (not review comments).

@github-actions

Copy link
Copy Markdown
Contributor

Reminder, please take a look at this pr: @shunping

@shunping

Copy link
Copy Markdown
Collaborator

r: @claudevdm

Could help to review this since you are more familiar with BigQueryIO? Thanks!

@shunping
shunping requested a review from claudevdm July 31, 2026 01:23
@github-actions

Copy link
Copy Markdown
Contributor

Stopping reviewer notifications for this pull request: review requested by someone other than the bot, ceding control. If you'd like to restart, comment assign set of reviewers

@claudevdm claudevdm left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

does not get a mapping of <destination, BeamSchema>

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

destination aware, trimming the TableRows to remove padded fields before writing etc.

The complexity and potential issues come from requiring a union schema.

Can we consider not forcing every row through one PCollection? How about we

  1. Group destinations by distinct schema in python
  2. Route each element to a branch matching its destination schema. Each branch has its own write path with a single concrete schema
  3. 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
  4. 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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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):

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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))

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can we also validate the expected schema of the tables that were created after the pipeline ran?

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants