Skip to content

Commit 4ac6c18

Browse files
committed
reenable window expiration context
1 parent 8ccf744 commit 4ac6c18

2 files changed

Lines changed: 4 additions & 7 deletions

File tree

sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignatures.java

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -222,7 +222,6 @@ private DoFnSignatures() {}
222222
private static final Collection<Class<? extends Parameter>>
223223
ALLOWED_ON_WINDOW_EXPIRATION_PARAMETERS =
224224
ImmutableList.of(
225-
Parameter.OnWindowExpirationContextParameter.class,
226225
Parameter.WindowParameter.class,
227226
Parameter.PipelineOptionsParameter.class,
228227
Parameter.OutputReceiverParameter.class,

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/SchemaUpdateHoldingFn.java

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -191,18 +191,16 @@ public void onPollTimer(
191191

192192
@OnWindowExpiration
193193
public void onWindowExpiration(
194-
// OnWindowExpirationContext context,
194+
OnWindowExpirationContext context,
195195
@Key ShardedKey<DestinationT> key,
196196
PipelineOptions pipelineOptions,
197197
@StateId("bufferedElements") BagState<TimestampedValue<ElementT>> bag,
198198
@StateId("minBufferedTimestamp") CombiningState<Long, long[], Long> minBufferedTimestamp,
199199
MultiOutputReceiver o)
200200
throws Exception {
201-
// TODO: Beam doesn't current support OnWindowExpirationContext. Reenable after this is fixed.
202-
// https://github.com/apache/beam/issues/38875
203-
// convertMessagesDoFn
204-
// .getDynamicDestinations()
205-
// .setSideInputAccessorFromOnWindowExpirationContext(context);
201+
convertMessagesDoFn
202+
.getDynamicDestinations()
203+
.setSideInputAccessorFromOnWindowExpirationContext(context);
206204

207205
// This can happen on test completion or drain. We can't set any more timers in window
208206
// expiration, so we just have to loop until the schema is updated.

0 commit comments

Comments
 (0)