[Dataflow Streaming Java] Fix possible IllegalStateException when grpc streams have deadline exceeded. - #36170
Conversation
…c streams have deadline exceeded.
There was a problem hiding this comment.
Summary of Changes
Hello @scwhittle, 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 resolves a critical IllegalStateException that could occur in Dataflow Streaming Java jobs when gRPC streams encountered deadline exceeded errors. The core issue stemmed from improper state management within ResettableThrowingStreamObserver and DirectStreamObserver, leading to onComplete being called after onError had already been processed. The changes introduce more robust state tracking and error handling logic, ensuring that stream observers correctly manage their lifecycle and prevent redundant or out-of-order termination calls, thereby improving the stability of streaming data processing.
Highlights
- Stream Observer State Management: Modified
ResettableThrowingStreamObserverto preventIllegalStateExceptionby ensuringisCurrentStreamClosedis correctly updated andonErroris not invoked on an already closed stream, particularly when handlingStreamObserverCancelledException. - Robust Stream Closure: Wrapped
onErrorandonCompletedcalls inResettableThrowingStreamObserverwithtry-finallyblocks to guarantee that theisCurrentStreamClosedflag is always set, even if the delegate observer throws an exception. - Outbound Observer Error Sequencing: Adjusted the
onErrormethod inDirectStreamObserverto set theisOutboundObserverClosedflag before calling theoutboundObserver.onErrormethod, preventing potential race conditions or re-entry issues. - Exception Visibility: Made the constructors for
StreamObserverCancelledExceptionpublic, likely to facilitate broader usage or testing of this specific exception type. - Enhanced Test Coverage: Added new unit tests to
ResettableThrowingStreamObserverTestto cover basiconNext,onError, andonCompletedscenarios, as well as specific cases involvingStreamObserverCancelledExceptionduringonNextprocessing.
Using Gemini Code Assist
The 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 /gemini <command> or @gemini-code-assist <command>. Below is a summary of the supported commands on the current page.
| Feature | Command | Description |
|---|---|---|
| Code Review | /gemini review |
Performs a code review for the current pull request in its current state. |
| Pull Request Summary | /gemini summary |
Provides a summary of the current pull request in its current state. |
| Comment | @gemini-code-assist | Responds in comments when explicitly tagged, both in pull request comments and review comments. |
| Help | /gemini help |
Displays a list of available commands. |
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 .gemini/ folder in the base of the repository. Detailed instructions can be found here.
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
-
Review the Privacy Notices, Generative AI Prohibited Use Policy, Terms of Service, and learn how to configure Gemini Code Assist in GitHub here. Gemini can make mistakes, so double check it and use code with caution. ↩
|
Stopping reviewer notifications for this pull request: review requested by someone other than the bot, ceding control. If you'd like to restart, comment |
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request effectively addresses a potential IllegalStateException in gRPC streams by improving the state management in ResettableThrowingStreamObserver under concurrent access. The core fix correctly handles stream closure when onNext fails, preventing race conditions with subsequent onCompleted or onError calls. The use of try-finally blocks in onError and onCompleted and the state update ordering change in DirectStreamObserver are excellent defensive programming practices that enhance robustness. The new tests provide good coverage for the fix. I have one suggestion to refactor the new tests for better maintainability.
…ache/beam/runners/dataflow/worker/windmill/client/ResettableThrowingStreamObserverTest.java Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>
| // Either the no longer active observer which we attempt a reset and handle errors | ||
| // or the current observer that still requires closing. | ||
| try { | ||
| delegate.onError(cancellationException); |
There was a problem hiding this comment.
this is under the mutex now, is it expected?
There was a problem hiding this comment.
I was doing it to be consistent with onError synchronization but I will change back since if we want to CP then it seems a riskier change.
|
Run Java PreCommit |
|
Not sure if related but GrpcGetDataStreamTest appears to be increasingly flaky: #36347 also observed in Beam import and repro'd locally update: The symptom remains the same after reverting this PR locally and run tests |
Example error that could be observed previously when onComplete was called on ResettableThrowingStreamObserver which already called delegate.onError within onNext.
{ "jsonPayload": { "exception": "java.lang.IllegalStateException\n at com.google.common.base.Preconditions.checkState(Preconditions.java:527)\n at org.apache.beam.runners.dataflow.worker.windmill.client.grpc.observers.DirectStreamObserver.onCompleted(DirectStreamObserver.java:198)\n at org.apache.beam.runners.dataflow.worker.windmill.client.ResettableThrowingStreamObserver.onCompleted(ResettableThrowingStreamObserver.java:143)\n at org.apache.beam.runners.dataflow.worker.windmill.client.AbstractWindmillStream.halfClose(AbstractWindmillStream.java:450)\n at org.apache.beam.runners.dataflow.worker.windmill.client.WindmillStreamPool.getStream(WindmillStreamPool.java:131)\n at org.apache.beam.runners.dataflow.worker.windmill.client.WindmillStreamPool.getCloseableStream(WindmillStreamPool.java:137)\n at org.apache.beam.runners.dataflow.worker.windmill.client.commits.StreamingEngineWorkCommitter.streamingCommitLoop(StreamingEngineWorkCommitter.java:178)\n at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1144)\n at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:642)\n at java.base/java.lang.Thread.run(Thread.java:1583)", "job": "2025-09-12_13_56_55-8555386915480239630", "logger": "org.apache.beam.runners.dataflow.worker.windmill.client.grpc.GrpcCommitWorkStream", "message": "Unexpected error when trying to close stream", "thread": "49", } }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.