goutamadwant opened a new pull request, #40294: URL: https://github.com/apache/beam/pull/40294
Addresses the remaining Java pending-result scalability item in #18459. `Watch` currently keeps every distinct unseen output from a poll in checkpointed pending restriction state. A very large poll can therefore create an oversized state value. ### What changes - Adds the opt-in `Watch.Growth.withMaxPendingResults(int)` API. - For finite, repeatable incomplete polls, keeps at most the configured number of oldest unseen outputs as pending work in each round. - Holds the watermark at or before the earliest omitted output and continues polling until omitted work can be recovered. - Defers a configured termination condition while a poll is truncated, avoiding silent loss of omitted outputs. - Preserves the existing first-output behavior for duplicate output keys. - Never limits `PollResult.complete()`, because its contract promises that the polling function will not be called again. ### Compatibility and scope The option is disabled by default, so existing pipelines retain their current behavior. It does not change the restriction coder or persisted state format. When enabled, the `PollFn` must repeat omitted outputs in subsequent polls and must eventually return a finite set that can drain below the limit. A function that continuously returns more distinct unseen outputs than the limit can keep polling indefinitely; the API documentation and runtime logging make this behavior explicit. The option bounds checkpointed pending outputs. It does not bound the list constructed by the polling function, transient hashes used during selection, or completed keys sharing a retained timestamp. It can be combined with `withTimestampCursor()` to time-bound completed-key state when output timestamps advance past the retention floor. This PR covers the Java SDK. Equivalent Python behavior can be tracked separately. ### Validation - `:sdks:java:core:test --tests 'org.apache.beam.sdk.transforms.WatchTest' --no-build-cache --rerun-tasks` - 30 tests passed - `:runners:direct-java:needsRunnerTests --tests 'org.apache.beam.sdk.transforms.WatchTest.testManyResultsWithPendingLimit' --no-build-cache --rerun-tasks` - passed; recovered 300 outputs over three bounded rounds - `:sdks:java:core:checkstyleMain` - `:sdks:java:core:checkstyleTest` - `:sdks:java:core:spotlessJavaCheck` - `:sdks:java:core:javadoc` - `:sdks:java:core:spotbugsMain` ------------------------ Thank you for your contribution! Follow this checklist to help us incorporate your contribution quickly and easily: - [x] Mention the appropriate issue in your description. - [x] Update `CHANGES.md` with noteworthy changes. - [ ] If this contribution is large, file an Apache [Individual Contributor License Agreement](https://www.apache.org/licenses/icla.pdf). See the [Contributor Guide](https://beam.apache.org/contribute) for more tips on [making the review process smoother](https://github.com/apache/beam/blob/master/CONTRIBUTING.md#make-the-reviewers-job-easier). To check the build health, see [BUILD_STATUS.md](https://github.com/apache/beam/blob/master/.test-infra/BUILD_STATUS.md). See [CI.md](https://github.com/apache/beam/blob/master/CI.md) and the [workflows README](https://github.com/apache/beam/blob/master/.github/workflows/README.md) for GitHub Actions details and workflow triggers. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
