From 7938f7f5aaca3ed85d11e10f6211c769733fff60 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Wed, 5 Aug 2026 23:43:20 +0000 Subject: [PATCH] [Dataflow Streaming] Remove finalizeCommits from processWork --- .../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 9e8265e509af..05a9ad82f182 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 @@ private void processWork( 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;