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());
}