This is an automated email from the ASF dual-hosted git repository.

scwhittle 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 eea1e03cf8e [Dataflow Streaming] Remove redundant onKeyTransition call 
(#39652)
eea1e03cf8e is described below

commit eea1e03cf8e9b30a26ae8ae19088ea6ba331f451
Author: Arun Pandian <[email protected]>
AuthorDate: Thu Aug 6 05:02:22 2026 -0700

    [Dataflow Streaming] Remove redundant onKeyTransition call (#39652)
---
 .../beam/runners/dataflow/worker/StreamingModeExecutionContext.java | 1 -
 .../runners/dataflow/worker/StreamingModeExecutionContextTest.java  | 6 ++++--
 2 files changed, 4 insertions(+), 3 deletions(-)

diff --git 
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java
 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java
index d577b861407..68dbd61f15f 100644
--- 
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java
+++ 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java
@@ -784,7 +784,6 @@ public class StreamingModeExecutionContext
       flushStateInternal();
       Work newWork = additionalWork.work();
       ++workItemsPolled;
-      checkStateNotNull(keyTransitionListener).onKeyTransition(activeWork, 
newWork);
       startForNewKey(newWork);
       return true;
     }
diff --git 
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java
 
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java
index c5efcea4e47..eb6bb51e420 100644
--- 
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java
+++ 
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java
@@ -550,11 +550,13 @@ public class StreamingModeExecutionContextTest {
         .thenReturn(executableWork2)
         .thenReturn(null);
 
-    executionContext.start(
-        work1, workExecutor, mockExecutor, mockHandle, null, (oldWork, 
newWork) -> {});
+    StreamingModeExecutionContext.KeyTransitionListener mockListener =
+        mock(StreamingModeExecutionContext.KeyTransitionListener.class);
+    executionContext.start(work1, workExecutor, mockExecutor, mockHandle, 
null, mockListener);
 
     assertTrue(executionContext.advance());
     assertEquals("key2", executionContext.getSerializedKey().toStringUtf8());
+    verify(mockListener, times(1)).onKeyTransition(work1, work2);
     assertFalse(executionContext.advance());
   }
 

Reply via email to