Skip to content

Commit c3ca6b5

Browse files
committed
foo
1 parent 4d82e8e commit c3ca6b5

7 files changed

Lines changed: 80 additions & 2 deletions

File tree

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

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

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -194,6 +194,7 @@ public void process(
194194

195195
@OnTimer("retryMismatchedRowsTimer")
196196
public void onTimer(
197+
OnTimerContext context,
197198
@Key ShardedKey<DestinationT> shardedDestination,
198199
@Timestamp Instant timestamp,
199200
@StateId("mismatchedRows") BagState<MismatchedRow> mismatchedRowsBag,
@@ -203,7 +204,7 @@ public void onTimer(
203204
PipelineOptions pipelineOptions,
204205
MultiOutputReceiver o)
205206
throws Exception {
206-
// TODO XXX SET DYNAMIC DESTINATIONS
207+
dynamicDestinations.setSideInputAccessorFromOnTimerContext(context);
207208

208209
mismatchedRowsBag.readLater();
209210
currentTimerValue.readLater();

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

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,33 @@ interface SideInputAccessor {
8282
private transient @Nullable SideInputAccessor sideInputAccessor;
8383
private transient @Nullable PipelineOptions options;
8484

85+
static class SideInputAccessorViaOnTimerContext implements SideInputAccessor {
86+
private DoFn<?, ?>.OnTimerContext onTimerContext;
87+
88+
public SideInputAccessorViaOnTimerContext(DoFn<?, ?>.OnTimerContext onTimerContext) {
89+
this.onTimerContext = onTimerContext;
90+
}
91+
92+
@Override
93+
public <SideInputT> SideInputT sideInput(PCollectionView<SideInputT> view) {
94+
return onTimerContext.sideInput(view);
95+
}
96+
}
97+
98+
static class SideInputAccessorViaOnWindowExpirationContext implements SideInputAccessor {
99+
private DoFn<?, ?>.OnWindowExpirationContext onWindowExpirationContext;
100+
101+
public SideInputAccessorViaOnWindowExpirationContext(
102+
DoFn<?, ?>.OnWindowExpirationContext onWindowExpirationContext) {
103+
this.onWindowExpirationContext = onWindowExpirationContext;
104+
}
105+
106+
@Override
107+
public <SideInputT> SideInputT sideInput(PCollectionView<SideInputT> view) {
108+
return onWindowExpirationContext.sideInput(view);
109+
}
110+
}
111+
85112
static class SideInputAccessorViaProcessContext implements SideInputAccessor {
86113
private DoFn<?, ?>.ProcessContext processContext;
87114

@@ -129,6 +156,17 @@ void setSideInputAccessorFromProcessContext(DoFn<?, ?>.ProcessContext context) {
129156
this.options = context.getPipelineOptions();
130157
}
131158

159+
void setSideInputAccessorFromOnTimerContext(DoFn<?, ?>.OnTimerContext context) {
160+
this.sideInputAccessor = new SideInputAccessorViaOnTimerContext(context);
161+
this.options = context.getPipelineOptions();
162+
}
163+
164+
void setSideInputAccessorFromOnWindowExpirationContext(
165+
DoFn<?, ?>.OnWindowExpirationContext context) {
166+
this.sideInputAccessor = new SideInputAccessorViaOnWindowExpirationContext(context);
167+
this.options = context.getPipelineOptions();
168+
}
169+
132170
/**
133171
* Returns an object that represents at a high level which table is being written to. May not
134172
* return null.

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -215,6 +215,19 @@ void setSideInputAccessorFromProcessContext(DoFn<?, ?>.ProcessContext context) {
215215
inner.setSideInputAccessorFromProcessContext(context);
216216
}
217217

218+
@Override
219+
void setSideInputAccessorFromOnTimerContext(DoFn<?, ?>.OnTimerContext context) {
220+
super.setSideInputAccessorFromOnTimerContext(context);
221+
inner.setSideInputAccessorFromOnTimerContext(context);
222+
}
223+
224+
@Override
225+
void setSideInputAccessorFromOnWindowExpirationContext(
226+
DoFn<?, ?>.OnWindowExpirationContext context) {
227+
super.setSideInputAccessorFromOnWindowExpirationContext(context);
228+
inner.setSideInputAccessorFromOnWindowExpirationContext(context);
229+
}
230+
218231
@Override
219232
public String toString() {
220233
return MoreObjects.toStringHelper(this).add("inner", inner).toString();

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

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -165,6 +165,7 @@ public Duration getAllowedTimestampSkew() {
165165

166166
@OnTimer("pollTimer")
167167
public void onPollTimer(
168+
OnTimerContext context,
168169
@Key ShardedKey<DestinationT> key,
169170
PipelineOptions pipelineOptions,
170171
@StateId("bufferedElements") BagState<TimestampedValue<ElementT>> bag,
@@ -174,6 +175,8 @@ public void onPollTimer(
174175
BoundedWindow window,
175176
MultiOutputReceiver o)
176177
throws Exception {
178+
convertMessagesDoFn.getDynamicDestinations().setSideInputAccessorFromOnTimerContext(context);
179+
177180
if (tryFlushBuffer(key.getKey(), pipelineOptions, bag, minBufferedTimestamp, o)) {
178181
timerTs.clear();
179182
} else {
@@ -188,12 +191,19 @@ public void onPollTimer(
188191

189192
@OnWindowExpiration
190193
public void onWindowExpiration(
194+
// OnWindowExpirationContext context,
191195
@Key ShardedKey<DestinationT> key,
192196
PipelineOptions pipelineOptions,
193197
@StateId("bufferedElements") BagState<TimestampedValue<ElementT>> bag,
194198
@StateId("minBufferedTimestamp") CombiningState<Long, long[], Long> minBufferedTimestamp,
195199
MultiOutputReceiver o)
196200
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);
206+
197207
// This can happen on test completion or drain. We can't set any more timers in window
198208
// expiration, so we just have to loop until the schema is updated.
199209
BackOff backoff =

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,4 +61,17 @@ void setSideInputAccessorFromProcessContext(DoFn<?, ?>.ProcessContext context) {
6161
super.setSideInputAccessorFromProcessContext(context);
6262
inner.setSideInputAccessorFromProcessContext(context);
6363
}
64+
65+
@Override
66+
void setSideInputAccessorFromOnTimerContext(DoFn<?, ?>.OnTimerContext context) {
67+
super.setSideInputAccessorFromOnTimerContext(context);
68+
inner.setSideInputAccessorFromOnTimerContext(context);
69+
}
70+
71+
@Override
72+
void setSideInputAccessorFromOnWindowExpirationContext(
73+
DoFn<?, ?>.OnWindowExpirationContext context) {
74+
super.setSideInputAccessorFromOnWindowExpirationContext(context);
75+
inner.setSideInputAccessorFromOnWindowExpirationContext(context);
76+
}
6477
}

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1238,6 +1238,7 @@ private void processPayloads(
12381238

12391239
@OnTimer("retryMismatchedRowsTimer")
12401240
public void onMismatchedRowsTimer(
1241+
OnTimerContext context,
12411242
PipelineOptions pipelineOptions,
12421243
@Key ShardedKey<DestinationT> shardedDestination,
12431244
@Timestamp org.joda.time.Instant elementTs,
@@ -1251,11 +1252,12 @@ public void onMismatchedRowsTimer(
12511252
@StateId("currentMismatchedRowTimerValue") ValueState<Long> currentTimerValue,
12521253
@StateId("minPendingTimestamp") ValueState<Long> minPendingTimestamp)
12531254
throws Exception {
1255+
dynamicDestinations.setSideInputAccessorFromOnTimerContext(context);
1256+
12541257
mismatchedRowsBag.readLater();
12551258
currentTimerValue.readLater();
12561259
minPendingTimestamp.readLater();
12571260

1258-
// TODOTDO XXX add context for side inputs.
12591261
TableDestination tableDestination =
12601262
destinations.computeIfAbsent(
12611263
shardedDestination.getKey(),

0 commit comments

Comments
 (0)