Skip to content

Commit eda08d8

Browse files
authored
Changes SplittableDoFn to call TruncateRestriction on drain (#39535)
* Changes SplittableDoFn to call TruncateRestriction first when getting a timer caused by drain. We then pass the residual restriction (if present) to ProcessElement. * Adds the rest of a comment string.
1 parent 6e2044b commit eda08d8

2 files changed

Lines changed: 333 additions & 20 deletions

File tree

runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java

Lines changed: 68 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -469,8 +469,9 @@ public String getErrorContext() {
469469
restrictionState.readLater();
470470
watermarkEstimatorState.readLater();
471471
WindowedValue<InputT> read = elementState.read();
472+
RestrictionT restriction = restrictionState.read();
472473
if (timer.causedByDrain() == CausedByDrain.CAUSED_BY_DRAIN) {
473-
read =
474+
WindowedValue<InputT> drainRead =
474475
WindowedValues.of(
475476
read.getValue(),
476477
read.getTimestamp(),
@@ -481,8 +482,73 @@ public String getErrorContext() {
481482
CausedByDrain.CAUSED_BY_DRAIN,
482483
read.getOpenTelemetryContext(),
483484
read.getValueKind());
485+
RestrictionTracker.TruncateResult<RestrictionT> truncateResult =
486+
invoker.invokeTruncateRestriction(
487+
new BaseArgumentProvider<InputT, OutputT>() {
488+
@Override
489+
public InputT element(DoFn<InputT, OutputT> doFn) {
490+
return drainRead.getValue();
491+
}
492+
493+
@Override
494+
public Object restriction() {
495+
return restriction;
496+
}
497+
498+
@Override
499+
public RestrictionTracker<?, ?> restrictionTracker() {
500+
return invoker.invokeNewTracker(this);
501+
}
502+
503+
@Override
504+
public Instant timestamp(DoFn<InputT, OutputT> doFn) {
505+
return drainRead.getTimestamp();
506+
}
507+
508+
@Override
509+
public PipelineOptions pipelineOptions() {
510+
return c.getPipelineOptions();
511+
}
512+
513+
@Override
514+
public PaneInfo paneInfo(DoFn<InputT, OutputT> doFn) {
515+
return drainRead.getPaneInfo();
516+
}
517+
518+
@Override
519+
public BoundedWindow window() {
520+
return Iterables.getOnlyElement(drainRead.getWindows());
521+
}
522+
523+
@Override
524+
public Object sideInput(String tagId) {
525+
PCollectionView<?> view = sideInputMapping.get(tagId);
526+
if (view == null) {
527+
throw new IllegalArgumentException(
528+
"calling getSideInput() with unknown view");
529+
}
530+
return sideInputReader.get(
531+
view, view.getWindowMappingFn().getSideInputWindow(window()));
532+
}
533+
534+
@Override
535+
public String getErrorContext() {
536+
return ProcessFn.class.getSimpleName() + ".invokeTruncateRestriction";
537+
}
538+
});
539+
if (truncateResult == null) {
540+
elementState.clear();
541+
restrictionState.clear();
542+
watermarkEstimatorState.clear();
543+
holdState.clear();
544+
return;
545+
}
546+
RestrictionT truncatedRestriction = truncateResult.getTruncatedRestriction();
547+
elementAndRestriction = KV.of(drainRead, truncatedRestriction);
548+
restrictionState.write(truncatedRestriction);
549+
} else {
550+
elementAndRestriction = KV.of(read, restriction);
484551
}
485-
elementAndRestriction = KV.of(read, restrictionState.read());
486552
watermarkEstimatorStateT = watermarkEstimatorState.read();
487553
}
488554

0 commit comments

Comments
 (0)