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]

Reply via email to