diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java index 470e22a66991..1825b77b65fb 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java @@ -649,6 +649,13 @@ public DoFn.OnTimerContext onTimerContext(DoFn "Cannot access OnTimerContext outside of @OnTimer methods."); } + @Override + public DoFn.OnWindowExpirationContext onWindowExpirationContext( + DoFn doFn) { + throw new UnsupportedOperationException( + "Cannot access OnWindowExpirationContext outside of @OnWindowExpiration methods."); + } + @Override public RestrictionTracker restrictionTracker() { throw new UnsupportedOperationException("RestrictionTracker parameters are not supported."); @@ -958,6 +965,13 @@ public DoFn.OnTimerContext onTimerContext(DoFn return this; } + @Override + public DoFn.OnWindowExpirationContext onWindowExpirationContext( + DoFn doFn) { + throw new UnsupportedOperationException( + "Cannot access OnWindowExpirationContext outside of @OnWindowExpiration methods."); + } + @Override public RestrictionTracker restrictionTracker() { throw new UnsupportedOperationException("RestrictionTracker parameters are not supported."); @@ -1299,6 +1313,12 @@ public DoFn.OnTimerContext onTimerContext(DoFn throw new UnsupportedOperationException("OnTimerContext parameters are not supported."); } + @Override + public DoFn.OnWindowExpirationContext onWindowExpirationContext( + DoFn doFn) { + return this; + } + @Override public RestrictionTracker restrictionTracker() { throw new UnsupportedOperationException("RestrictionTracker parameters are not supported."); diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/ByteBuddyDoFnInvokerFactory.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/ByteBuddyDoFnInvokerFactory.java index 3ebabb6e3c37..c08243fda5c0 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/ByteBuddyDoFnInvokerFactory.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/ByteBuddyDoFnInvokerFactory.java @@ -142,6 +142,8 @@ class ByteBuddyDoFnInvokerFactory implements DoFnInvokerFactory { public static final String OUTPUT_PARAMETER_METHOD = "outputReceiver"; public static final String TAGGED_OUTPUT_PARAMETER_METHOD = "taggedOutputReceiver"; public static final String ON_TIMER_CONTEXT_PARAMETER_METHOD = "onTimerContext"; + public static final String ON_WINDOW_EXPIRATION_CONTEXT_PARAMETER_METHOD = + "onWindowExpirationContext"; public static final String WINDOW_PARAMETER_METHOD = "window"; public static final String PANE_INFO_PARAMETER_METHOD = "paneInfo"; public static final String PIPELINE_OPTIONS_PARAMETER_METHOD = "pipelineOptions"; @@ -1170,6 +1172,16 @@ public StackManipulation dispatch(OnTimerContextParameter p) { ON_TIMER_CONTEXT_PARAMETER_METHOD, DoFn.class))); } + @Override + public StackManipulation dispatch( + DoFnSignature.Parameter.OnWindowExpirationContextParameter p) { + return new StackManipulation.Compound( + pushDelegate, + MethodInvocation.invoke( + getExtraContextFactoryMethodDescription( + ON_WINDOW_EXPIRATION_CONTEXT_PARAMETER_METHOD, DoFn.class))); + } + @Override public StackManipulation dispatch(WindowParameter p) { return new StackManipulation.Compound( diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnInvoker.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnInvoker.java index eaabdff907c7..c8c7ddf24b69 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnInvoker.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnInvoker.java @@ -185,6 +185,10 @@ interface ArgumentProvider { /** Provide a {@link DoFn.OnTimerContext} to use with the given {@link DoFn}. */ DoFn.OnTimerContext onTimerContext(DoFn doFn); + /** Provide a {@link DoFn.OnWindowExpirationContext} to use with the given {@link DoFn}. */ + DoFn.OnWindowExpirationContext onWindowExpirationContext( + DoFn doFn); + /** Provide a reference to the input element. */ InputT element(DoFn doFn); @@ -447,6 +451,13 @@ public DoFn.OnTimerContext onTimerContext(DoFn String.format("OnTimerContext unsupported in %s", getErrorContext())); } + @Override + public DoFn.OnWindowExpirationContext onWindowExpirationContext( + DoFn doFn) { + throw new UnsupportedOperationException( + String.format("OnWindowExpirationContext unsupported in %s", getErrorContext())); + } + @Override public State state(String stateId, boolean alwaysFetched) { throw new UnsupportedOperationException( @@ -538,6 +549,12 @@ public DoFn.OnTimerContext onTimerContext(DoFn return delegate.onTimerContext(doFn); } + @Override + public DoFn.OnWindowExpirationContext onWindowExpirationContext( + DoFn doFn) { + return delegate.onWindowExpirationContext(doFn); + } + @Override public InputT element(DoFn doFn) { return delegate.element(doFn); diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignature.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignature.java index 51dadd178a6f..99b002c1106d 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignature.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignature.java @@ -305,6 +305,8 @@ public ResultT match(Cases cases) { return cases.dispatch((ProcessContextParameter) this); } else if (this instanceof OnTimerContextParameter) { return cases.dispatch((OnTimerContextParameter) this); + } else if (this instanceof OnWindowExpirationContextParameter) { + return cases.dispatch((OnWindowExpirationContextParameter) this); } else if (this instanceof WindowParameter) { return cases.dispatch((WindowParameter) this); } else if (this instanceof PaneInfoParameter) { @@ -391,6 +393,8 @@ public interface Cases { ResultT dispatch(OnTimerContextParameter p); + ResultT dispatch(OnWindowExpirationContextParameter p); + ResultT dispatch(WindowParameter p); ResultT dispatch(PaneInfoParameter p); @@ -498,6 +502,11 @@ public ResultT dispatch(OnTimerContextParameter p) { return dispatchDefault(p); } + @Override + public ResultT dispatch(OnWindowExpirationContextParameter p) { + return dispatchDefault(p); + } + @Override public ResultT dispatch(WindowParameter p) { return dispatchDefault(p); diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignatures.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignatures.java index 2983fc94021c..9f3491bca7b9 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignatures.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignatures.java @@ -229,7 +229,8 @@ private DoFnSignatures() {} Parameter.StateParameter.class, Parameter.TimestampParameter.class, Parameter.KeyParameter.class, - Parameter.SideInputParameter.class); + Parameter.SideInputParameter.class, + Parameter.OnWindowExpirationContextParameter.class); private static final Collection> ALLOWED_GET_INITIAL_RESTRICTION_PARAMETERS = diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/SplittableParDoNaiveBounded.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/SplittableParDoNaiveBounded.java index d1fb23e77c47..52b1174e3a03 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/SplittableParDoNaiveBounded.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/SplittableParDoNaiveBounded.java @@ -506,6 +506,12 @@ public DoFn.OnTimerContext onTimerContext(DoFn throw new IllegalStateException(); } + @Override + public DoFn.OnWindowExpirationContext onWindowExpirationContext( + DoFn doFn) { + throw new IllegalStateException(); + } + @Override public InputT element(DoFn doFn) { return element; diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ParDoTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ParDoTest.java index 0c984d01c8f0..6beea338689b 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ParDoTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ParDoTest.java @@ -7335,11 +7335,15 @@ public void onTimer( @OnWindowExpiration public void onWindowExpiration( @AlwaysFetched @StateId(stateId) ValueState state, + BoundedWindow window, @Key String key, + OnWindowExpirationContext context, OutputReceiver r) { Integer currentValue = MoreObjects.firstNonNull(state.read(), 0); // verify state assertEquals(1, (int) currentValue); + Preconditions.checkNotNull(context); + assertEquals(window, context.window()); // To check output is received from OnWindowExpiration r.output(currentValue); } diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/reflect/DoFnSignaturesTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/reflect/DoFnSignaturesTest.java index 5a5353482c95..330bb5b94418 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/reflect/DoFnSignaturesTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/reflect/DoFnSignaturesTest.java @@ -1422,16 +1422,20 @@ public void bar( @StateId("foo") ValueState s, PipelineOptions p, OutputReceiver o, - MultiOutputReceiver m) {} + MultiOutputReceiver m, + OnWindowExpirationContext c) {} }.getClass()); List params = sig.onWindowExpiration().extraParameters(); - assertThat(params.size(), equalTo(5)); + assertThat(params.size(), equalTo(6)); assertThat(params.get(0), instanceOf(WindowParameter.class)); assertThat(params.get(1), instanceOf(StateParameter.class)); assertThat(params.get(2), instanceOf(PipelineOptionsParameter.class)); assertThat(params.get(3), instanceOf(OutputReceiverParameter.class)); assertThat(params.get(4), instanceOf(TaggedOutputReceiverParameter.class)); + assertThat( + params.get(5), + instanceOf(DoFnSignature.Parameter.OnWindowExpirationContextParameter.class)); } private interface FeatureTest { diff --git a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnApiDoFnRunner.java b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnApiDoFnRunner.java index 3e4675ab074a..d5b5cebadb34 100644 --- a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnApiDoFnRunner.java +++ b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnApiDoFnRunner.java @@ -2250,6 +2250,13 @@ public DoFn.OnTimerContext onTimerContext(DoFn "Cannot access OnTimerContext outside of @OnTimer methods."); } + @Override + public DoFn.OnWindowExpirationContext onWindowExpirationContext( + DoFn doFn) { + throw new UnsupportedOperationException( + "Cannot access OnWindowExpirationContext outside of @OnWindowExpiration methods."); + } + @Override public RestrictionTracker restrictionTracker() { return currentTracker; @@ -2469,6 +2476,12 @@ private void checkOnWindowExpirationTimestamp(Instant timestamp) { private final OnWindowExpirationContext.Context context = new OnWindowExpirationContext.Context(); + @Override + public DoFn.OnWindowExpirationContext onWindowExpirationContext( + DoFn doFn) { + return context; + } + @Override public BoundedWindow window() { return currentWindow;