Fix output bugs in FnApiDoFnRunner @OnTimer and @OnWindowExpiration - #40024
Open
kennknowles wants to merge 1 commit into
Open
Fix output bugs in FnApiDoFnRunner @OnTimer and @OnWindowExpiration#40024kennknowles wants to merge 1 commit into
kennknowles wants to merge 1 commit into
Conversation
The three argument providers each carry their own near-identical copy of the output-receiver plumbing, and the @ontimer and @OnWindowExpiration copies had drifted from the working @ProcessElement one: - OnTimerContext.outputWindowedValue(tag, ...) had an empty body, so tagged windowed output from @ontimer was silently dropped. - OnTimerContext's tagged receivers passed no tag to outputWindowedValue, so output to a side tag went to the main output. - The @ontimer row receivers built from WindowedValues.builder( currentElement) or called withValue() on a fresh builder. Neither works during timer processing: currentElement is only set while processing an element, and withValue() reads timestamp/window/pane off the builder it is called on. Both now seed from currentTimer, matching the @OnWindowExpiration equivalents. - OnWindowExpirationContext read currentElement.getValueKind() while currentElement was null. Window expiration emits new records, so this is ValueKind.INSERT. - OnWindowExpirationContext.timeDomain() returned currentTimeDomain, which processOnWindowExpiration never sets. TimeDomain is not an allowed @OnWindowExpiration parameter, so drop the override and let BaseArgumentProvider reject it. - outputWindowedValue(tag, ...) skipped the unknown-tag check its siblings perform, so a bad tag produced an NPE. - processTimer's finally did not clear causedByDrain, unlike processOnWindowExpiration. Not observable today since every entry point sets it before invoking user code, but the asymmetry is a trap. The two inner Context classes are renamed to TimerContext and WindowExpirationContext so the fields holding them can drop their raw types; the enclosing providers are generic in K, and errorprone's SameNameButDifferent rejects the bare Context name. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
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.
The three argument providers each carry their own near-identical copy of the output-receiver plumbing, and the
@OnTimerand@OnWindowExpirationcopies had drifted from the working@ProcessElementone:OnTimerContext.outputWindowedValue(tag, ...)had an empty body, so tagged windowed output from@OnTimerwas silently dropped.OnTimerContext's tagged receivers passed no tag to outputWindowedValue, so output to a side tag went to the main output.@OnTimerrow receivers built fromWindowedValues.builder( currentElement)or calledwithValue()on a fresh builder. Neither works during timer processing: currentElement is only set while processing an element, and withValue() reads timestamp/window/pane off the builder it is called on. Both now seed from currentTimer, matching the@OnWindowExpirationequivalents.OnWindowExpirationContextread currentElement.getValueKind() while currentElement was null. Window expiration emits new records, so this isValueKind.INSERT.OnWindowExpirationContext.timeDomain()returnedcurrentTimeDomain, whichprocessOnWindowExpirationnever sets.TimeDomainis not an allowed@OnWindowExpiration parameter, so drop the override and letBaseArgumentProviderreject it.outputWindowedValue(tag, ...)skipped the unknown-tag check its siblings perform, so a bad tag produced an NPE.processTimer's finally did not clearcausedByDrain,unlike processOnWindowExpiration. Not observable today since every entry point sets it before invoking user code, but the asymmetry is a trap.The two inner Context classes are renamed to TimerContext and WindowExpirationContext so the fields holding them can drop their raw types; the enclosing providers are generic in K, and errorprone's SameNameButDifferent rejects the bare Context name.
Please add a meaningful description for your change here
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.