AddFiles: read side of the schema pre-pass - #39933
Open
claudevdm wants to merge 1 commit into
Open
Conversation
…emas.canonical, CollectDistinctSchemas) New, not yet wired code that turns a PCollection of file paths into the list of distinct schemas those files carry, with file counts. The pre-pass that will use it exists because manifest entries are immutable: a file registered before the table knows one of its columns never gets stats for that column, so the table schema has to be brought up to date before any file is registered, and that requires looking at every footer first. ReadFooterSchema (DoFn<String, String>) reads each Parquet footer on the BoundedAsyncTasks pool and emits the file's canonical schema as JSON. Non-Parquet paths and unknown extensions contribute nothing; a footer that cannot be read or converted is logged and counted (numFooterReadErrors) but never fails the pipeline: the per-file registration step reports such files individually later. FileSchemas.canonical sorts struct fields by name at every level and renumbers ids in deterministic order. The ids are positional and never consumed downstream: the commit side reconciles columns by name (unionByNameWith). CollectDistinctSchemas is a CombineFn over the canonical JSON strings (Map<String, Long> accumulator) producing List<KV<String, Long>> ordered most common first, ties broken by the JSON text for determinism. The most common schema goes first because the commit side uses it as the seed when the table does not exist yet.
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.
ReadFooterSchema, FileScemas.canonical, CollectDistinctSchemas)
New, not yet wired code that turns a PCollection of file paths into the list of distinct schemas those files carry, with file counts. The pre-pass that will use it exists because manifest entries are immutable: a file registered before the table knows one of its columns never gets stats for that column, so the table schema has to be brought up to date before any file is registered, and that requires looking at every footer first.
ReadFooterSchema (DoFn<String, String>) reads each Parquet footer on the BoundedAsyncTasks pool and emits the file's canonical schema as JSON. Non-Parquet paths and unknown extensions contribute nothing; a footer that cannot be read or converted is logged and counted (numFooterReadErrors) but never fails the pipeline: the per-file registration step reports such files individually later.
FileSchemas.canonical sorts struct fields by name at every level and renumbers ids in deterministic order.
The ids are positional and never consumed downstream: the commit side reconciles columns by name (unionByNameWith).
CollectDistinctSchemas is a CombineFn over the canonical JSON strings (Map<String, Long> accumulator) producing List<KV<String, Long>> ordered most common first, ties broken by the JSON text for determinism. The most common schema goes first because the commit side uses it as the seed when the table does not exist yet.
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.