Auto-infer BigQuery schema from schema'd PCollections for FILE_LOADS - #39662
Auto-infer BigQuery schema from schema'd PCollections for FILE_LOADS#39662nikitagrover19 wants to merge 4 commits into
Conversation
Adds beam_schema_to_bq_table_schema() to convert a Beam schema into a BigQuery TableSchema dict, and uses it in WriteToBigQuery's FILE_LOADS path (JSON and AVRO) when no explicit schema is provided but the input PCollection carries a Beam schema (e.g. NamedTuples, dataclasses, or Beam Rows). Mirrors the existing STORAGE_WRITE_API auto-inference. Fixes apache#28162
|
Checks are failing. Will not request review until checks are succeeding. If you'd like to override that behavior, comment |
|
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 |
|
Assigning new set of reviewers because Pr has gone too long without review. If you would like to opt out of this review, comment R: @jrmccluskey for label python. Available commands:
|
|
Reminder, please take a look at this pr: @jrmccluskey |
|
Assigning new set of reviewers because Pr has gone too long without review. If you would like to opt out of this review, comment R: @shunping for label python. Available commands:
|
Fixes #28162
What
WriteToBigQuerywithmethod=FILE_LOADSpreviously required an explicitschema, even when the input PCollection already carried a Beam schema - for example a PCollection produced byReadFromBigQuery(..., output_type='BEAM_ROW'), or any PCollection ofNamedTuples, dataclasses, or Beam Rows. The
STORAGE_WRITE_APImethodalready auto-infers a schema in this situation; this brings
FILE_LOADS(both JSON and AVRO temp file formats) in line with that behavior.
How
beam_schema_to_bq_table_schema()inbigquery_schema_tools.py,the reverse of the existing
generate_user_type_from_bq_schema(). Itconverts a Beam
schema_pb2.Schemainto a BigQueryTableSchemadict,handling atomic types, logical types (
Timestamp,Date,Decimal),nested/repeated fields, and nullability, and raises a clear
ValueErrorfor types with no BigQuery equivalent (e.g.
MapType).WriteToBigQuery.expand()'sFILE_LOADSbranch now triesschema_from_element_type(pcoll.element_type)whenself.schema is None, before the existing "schema required for AVRO" check, so AVROloads benefit from inference too, not just JSON.
schemastaysNone) when thePCollection has no attached Beam schema (e.g. plain dicts), so this is
purely additive and non-breaking.
Testing
decimals, nullable fields, repeated fields, nested records, repeated
records, and the unsupported-
MapTypeerror case(
bigquery_schema_tools_test.py).WriteToBigQueryconfirming: schemaauto-inference for JSON file loads, schema auto-inference for AVRO file
loads (previously always raised
"A schema must be provided"here),and that plain-dict PCollections are unaffected
(
bigquery_test.py).TestWriteToBigQuerytest class locally to check forregressions; all passing except one pre-existing test that requires
real GCP credentials unavailable in my local environment (unrelated to
this change).
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.