Skip to content

Commit 705db25

Browse files
authored
[Dataflow Streaming] Activate SourceState Finalizers before submitting workitem to harness threads (#38921)
1 parent a88686c commit 705db25

1 file changed

Lines changed: 2 additions & 3 deletions

File tree

  • runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -216,6 +216,8 @@ public void scheduleWork(
216216
Work.ProcessingContext processingContext,
217217
boolean drainMode,
218218
ImmutableList<LatencyAttribution> getWorkStreamLatencies) {
219+
// Before any processing starts, call any pending OnCommit callbacks
220+
commitFinalizer.finalizeCommits(workItem.getSourceState().getFinalizeIdsList());
219221
computationState.activateWork(
220222
ExecutableWork.create(
221223
Work.create(
@@ -255,9 +257,6 @@ private void processWork(
255257
setUpWorkLoggingContext(work.getLatencyTrackingId(), computationId);
256258
LOG.debug("Starting processing for {}:\n{}", computationId, work);
257259

258-
// Before any processing starts, call any pending OnCommit callbacks. Nothing that requires
259-
// cleanup should be done before this, since we might exit early here.
260-
commitFinalizer.finalizeCommits(workItem.getSourceState().getFinalizeIdsList());
261260
if (workItem.getSourceState().getOnlyFinalize()) {
262261
Windmill.WorkItemCommitRequest.Builder outputBuilder = initializeOutputBuilder(key, workItem);
263262
outputBuilder.setSourceStateUpdates(Windmill.SourceState.newBuilder().setOnlyFinalize(true));

0 commit comments

Comments
 (0)