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;

Reply via email to