-
Notifications
You must be signed in to change notification settings - Fork 4.4k
Enforce deterministic field order in Schema generated from KafkaIO classes. #36295
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Conversation
Summary of ChangesHello @baeminbo, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request enhances the robustness of Apache Beam's KafkaIO component by ensuring that schemas generated for its internal classes maintain a deterministic field order. By explicitly defining the field numbering through annotations, the change mitigates potential serialization inconsistencies and improves the reliability of data processing with KafkaIO. Additionally, new tests have been added to validate these schema generations. Highlights
Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counter productive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for Github and other Google products, sign up here. You can also get AI-powered code generation, chat, as well as code reviews directly in the IDE at no cost with the Gemini Code Assist IDE Extension. Footnotes
|
|
Checks are failing. Will not request review until checks are succeeding. If you'd like to override that behavior, comment |
…asses. To prevent potential serialization issues, this change adds the `@SchemaFieldNumber` annotation to KafkaIO classes that use `DefaultSchema`. This guarantees a deterministic field order, which is not ensured when schemas are generated with `AutoValueSchema`, `JavaBeanSchema`, or `JavaFieldSchema`.
61a18b5 to
15bb463
Compare
|
Run Java PreCommit |
|
Assigning reviewers: R: @ahmedabu98 for label java. 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: @ahmedabu98 @fozzie15 |
|
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: @Abacn for label java. Available commands:
|
|
@tomstepp @johnjcasey could take a look? |
|
LGTM in general. Assuming pipelines are update compatible across this change? @baeminbo |
Yes, right. Currently, this issue can randomly corrupt the states persisted in storage with |
|
To clarify, it sounds like this does fix updates where the original pipeline includes this change. But I was not sure if it works with original pipeline on an SDK version before this change upgrading to newer SDK version with this change. Would the original order be incompatible with the new enforced order? |
|
@johnjcasey or @an2x would you have any additional feedback? |
|
LGTM |
|
This breaks PythonUdingJavaSQL tests: test_streaming_taxifare_prediction_yaml (apache_beam.yaml.examples.testing.examples_test.MLTest) failed |
|
Seems like this also broke Managed I/O for Kafka (Google internal test) Failed to upgrade using the expansion service manager: INTERNAL: Expansion request failed: java.lang.IllegalStateException: Expected field number 3 for field + redistributeNumKeys instead got 2" Can we revert ? |
|
Do we know the Managed Kafka test expects specific field numbers? These
classes already had unstable field numbers that were liable to change -
this PR fixed that.
…On Mon, Oct 13, 2025 at 10:18 AM Chamikara Jayalath < ***@***.***> wrote:
*chamikaramj* left a comment (apache/beam#36295)
<#36295 (comment)>
Seems like this also broke Managed I/O for Kafka (Google internal test)
Failed to upgrade using the expansion service manager: INTERNAL: Expansion
request failed: java.lang.IllegalStateException: Expected field number 3
for field + redistributeNumKeys instead got 2"
Can we revert ?
—
Reply to this email directly, view it on GitHub
<#36295 (comment)>, or
unsubscribe
<https://github.com/notifications/unsubscribe-auth/AFAYJVJICPBLS72E5DH5LPT3XPNFDAVCNFSM6AAAAACHSAZFNGVHI2DSMVQWIX3LMV43OSLTON2WKQ3PNVWWK3TUHMZTGOJYGQYTAMZUG4>
.
You are receiving this because you modified the open/close state.Message
ID: ***@***.***>
|
|
Managed I/O uses a Beam construction Row (shema-aware transform) generated from old Beam versions to upgrade Kafka to newer Beam versions. So my suspicion is that that Row is now incompatible with the new Beam Kafka I/O version so we cannot upgrade. That test has been very stable historically in my experience (over last 10 Beam versions). May be we can add a guard against old Beam versions using the updateCompatibilityBeamVersion option [1] ? (which is always set for Managed I/O) beam/sdks/java/core/src/main/java/org/apache/beam/sdk/options/StreamingOptions.java Line 45 in 1c6f779
|
|
Full stack-trace: |
|
BTW we also sort fields when performing the schema-transform expansion: Line 155 in 1c6f779
cc: @ahmedabu98 |
|
Created a release blocker #36496 for tracking. |
|
The broken |
To prevent potential serialization issues, this change adds the
@SchemaFieldNumberannotation to KafkaIO classes that useDefaultSchema. This guarantees a deterministic field order, which is not ensured when schemas are generated withAutoValueSchema,JavaBeanSchema, orJavaFieldSchema.related to: #30276, b/443606305#comment29
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.