This is an automated email from the ASF dual-hosted git repository.
Abacn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 02ff2978792 [Dataflow Streaming] Remove finalizeCommits from
processWork (#39648)
02ff2978792 is described below
commit 02ff2978792bcdd17339312e5c488af5a9e17a2b
Author: Arun Pandian <[email protected]>
AuthorDate: Wed Aug 5 21:17:47 2026 -0700
[Dataflow Streaming] Remove finalizeCommits from processWork (#39648)
---
.../worker/windmill/work/processing/StreamingWorkScheduler.java | 4 ----
1 file changed, 4 deletions(-)
diff --git
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java
index 7c65c3326c9..299c67128ca 100644
---
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java
+++
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java
@@ -232,10 +232,6 @@ public class StreamingWorkScheduler {
KeyTransitionListener keyTransitionListener =
createKeyTransitionListener();
keyTransitionListener.onKeyTransition(null, work);
- // Before any processing starts, call any pending OnCommit callbacks.
Nothing that requires
- // cleanup should be done before this, since we might exit early here.
-
commitFinalizer.finalizeCommits(workItem.getSourceState().getFinalizeIdsList());
-
if (workItem.getSourceState().getOnlyFinalize()) {
handleOnlyFinalize(computationState, work, workItem);
return;