From d64384ea1f82d97e874132bc5c2bb93882ac81bd Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Thu, 6 Aug 2026 05:24:18 +0000 Subject: [PATCH] [Dataflow Streaming] Remove redundant onKeyTransition call --- .../dataflow/worker/StreamingModeExecutionContext.java | 1 - .../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 d577b8614078..68dbd61f15f6 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 boolean advance() throws CoderException { 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 c5efcea4e47c..eb6bb51e4207 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 void testAdvance_success() throws Exception { .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()); }