Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -240,6 +240,9 @@ public class DoFnOperator<PreInputT, InputT, OutputT>
/** Metrics container for reporting Beam metrics to Flink (null if metrics are disabled). */
transient @Nullable FlinkMetricContainer flinkMetricContainer;

/** Whether {@link #finish()} flushed all data, which only happens once the input has ended. */
private transient boolean finished;

/** Helper class to report the checkpoint duration. */
private transient @Nullable CheckpointStats checkpointStats;

Expand Down Expand Up @@ -299,7 +302,8 @@ public DoFnOperator(
this.sideInputTagMapping = sideInputTagMapping;
this.sideInputs = sideInputs;
this.serializedOptions = new SerializablePipelineOptions(options);
this.isStreaming = serializedOptions.get().as(FlinkPipelineOptions.class).isStreaming();
FlinkPipelineOptions flinkOptions = options.as(FlinkPipelineOptions.class);
this.isStreaming = flinkOptions.isStreaming();
this.windowingStrategy = windowingStrategy;
this.outputManagerFactory = outputManagerFactory;

Expand All @@ -311,8 +315,6 @@ public DoFnOperator(
this.timerCoder =
TimerInternals.TimerDataCoderV2.of(windowingStrategy.getWindowFn().windowCoder());

FlinkPipelineOptions flinkOptions = options.as(FlinkPipelineOptions.class);

this.maxBundleSize = flinkOptions.getMaxBundleSize();
Preconditions.checkArgument(maxBundleSize > 0, "Bundle size must be at least 1");
this.maxBundleTimeMills = flinkOptions.getMaxBundleTimeMills();
Expand Down Expand Up @@ -630,13 +632,26 @@ private void earlyBindStateIfNeeded() throws IllegalArgumentException, IllegalAc
}

void cleanUp() throws Exception {
Optional.ofNullable(flinkMetricContainer)
.ifPresent(FlinkMetricContainer::registerMetricsForPipelineResult);
if (shouldRegisterMetricsForPipelineResult()) {
Optional.ofNullable(flinkMetricContainer)
.ifPresent(FlinkMetricContainer::registerMetricsForPipelineResult);
}
Optional.ofNullable(checkFinishBundleTimer).ifPresent(timer -> timer.cancel(true));
Workarounds.deleteStaticCaches();
Optional.ofNullable(doFnInvoker).ifPresent(DoFnInvoker::invokeTeardown);
}

/**
* The metrics accumulator hands final totals to the program that launched the pipeline once the
* pipeline has ended. {@link #close()} also runs when the task is cancelled, which in streaming
* mode happens on every failover without the pipeline ending. Copying the metrics there is
* unnecessary and makes every JobManager status request merge and render the full metric state of
* each cancelled subtask. In streaming mode, only copy them once the input has ended.
*/
private boolean shouldRegisterMetricsForPipelineResult() {
return finished || !isStreaming;
}

void flushData() throws Exception {
// This is our last change to block shutdown of this operator while
// there are still remaining processing-time timers. Flink will ignore pending
Expand Down Expand Up @@ -689,6 +704,7 @@ void flushData() throws Exception {
public void finish() throws Exception {
try {
flushData();
finished = true;
} finally {
super.finish();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -239,6 +239,9 @@ public class DoFnOperator<PreInputT, InputT, OutputT>
/** Metrics container for reporting Beam metrics to Flink (null if metrics are disabled). */
transient @Nullable FlinkMetricContainer flinkMetricContainer;

/** Whether {@link #finish()} flushed all data, which only happens once the input has ended. */
private transient boolean finished;

/** Helper class to report the checkpoint duration. */
private transient @Nullable CheckpointStats checkpointStats;

Expand Down Expand Up @@ -298,7 +301,8 @@ public DoFnOperator(
this.sideInputTagMapping = sideInputTagMapping;
this.sideInputs = sideInputs;
this.serializedOptions = new SerializablePipelineOptions(options);
this.isStreaming = serializedOptions.get().as(FlinkPipelineOptions.class).isStreaming();
FlinkPipelineOptions flinkOptions = options.as(FlinkPipelineOptions.class);
this.isStreaming = flinkOptions.isStreaming();
this.windowingStrategy = windowingStrategy;
this.outputManagerFactory = outputManagerFactory;

Expand All @@ -311,8 +315,6 @@ public DoFnOperator(
this.timerCoder =
TimerInternals.TimerDataCoderV2.of(windowingStrategy.getWindowFn().windowCoder());

FlinkPipelineOptions flinkOptions = options.as(FlinkPipelineOptions.class);

this.maxBundleSize = flinkOptions.getMaxBundleSize();
Preconditions.checkArgument(maxBundleSize > 0, "Bundle size must be at least 1");
this.maxBundleTimeMills = flinkOptions.getMaxBundleTimeMills();
Expand Down Expand Up @@ -630,13 +632,26 @@ private void earlyBindStateIfNeeded() throws IllegalArgumentException, IllegalAc
}

void cleanUp() throws Exception {
Optional.ofNullable(flinkMetricContainer)
.ifPresent(FlinkMetricContainer::registerMetricsForPipelineResult);
if (shouldRegisterMetricsForPipelineResult()) {
Optional.ofNullable(flinkMetricContainer)
.ifPresent(FlinkMetricContainer::registerMetricsForPipelineResult);
}
Optional.ofNullable(checkFinishBundleTimer).ifPresent(timer -> timer.cancel(true));
Workarounds.deleteStaticCaches();
Optional.ofNullable(doFnInvoker).ifPresent(DoFnInvoker::invokeTeardown);
}

/**
* The metrics accumulator hands final totals to the program that launched the pipeline once the
* pipeline has ended. {@link #close()} also runs when the task is cancelled, which in streaming
* mode happens on every failover without the pipeline ending. Copying the metrics there is
* unnecessary and makes every JobManager status request merge and render the full metric state of
* each cancelled subtask. In streaming mode, only copy them once the input has ended.
*/
private boolean shouldRegisterMetricsForPipelineResult() {
return finished || !isStreaming;
}

void flushData() throws Exception {
// This is our last change to block shutdown of this operator while
// there are still remaining processing-time timers. Flink will ignore pending
Expand Down Expand Up @@ -689,6 +704,7 @@ void flushData() throws Exception {
public void finish() throws Exception {
try {
flushData();
finished = true;
} finally {
super.finish();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -239,6 +239,9 @@ public class DoFnOperator<PreInputT, InputT, OutputT>
/** Metrics container for reporting Beam metrics to Flink (null if metrics are disabled). */
transient @Nullable FlinkMetricContainer flinkMetricContainer;

/** Whether {@link #finish()} flushed all data, which only happens once the input has ended. */
private transient boolean finished;

/** Helper class to report the checkpoint duration. */
private transient @Nullable CheckpointStats checkpointStats;

Expand Down Expand Up @@ -298,7 +301,8 @@ public DoFnOperator(
this.sideInputTagMapping = sideInputTagMapping;
this.sideInputs = sideInputs;
this.serializedOptions = new SerializablePipelineOptions(options);
this.isStreaming = serializedOptions.get().as(FlinkPipelineOptions.class).isStreaming();
FlinkPipelineOptions flinkOptions = options.as(FlinkPipelineOptions.class);
this.isStreaming = flinkOptions.isStreaming();
this.windowingStrategy = windowingStrategy;
this.outputManagerFactory = outputManagerFactory;

Expand All @@ -311,8 +315,6 @@ public DoFnOperator(
this.timerCoder =
TimerInternals.TimerDataCoderV2.of(windowingStrategy.getWindowFn().windowCoder());

FlinkPipelineOptions flinkOptions = options.as(FlinkPipelineOptions.class);

this.maxBundleSize = flinkOptions.getMaxBundleSize();
Preconditions.checkArgument(maxBundleSize > 0, "Bundle size must be at least 1");
this.maxBundleTimeMills = flinkOptions.getMaxBundleTimeMills();
Expand Down Expand Up @@ -630,13 +632,26 @@ private void earlyBindStateIfNeeded() throws IllegalArgumentException, IllegalAc
}

void cleanUp() throws Exception {
Optional.ofNullable(flinkMetricContainer)
.ifPresent(FlinkMetricContainer::registerMetricsForPipelineResult);
if (shouldRegisterMetricsForPipelineResult()) {
Optional.ofNullable(flinkMetricContainer)
.ifPresent(FlinkMetricContainer::registerMetricsForPipelineResult);
}
Optional.ofNullable(checkFinishBundleTimer).ifPresent(timer -> timer.cancel(true));
Workarounds.deleteStaticCaches();
Optional.ofNullable(doFnInvoker).ifPresent(DoFnInvoker::invokeTeardown);
}

/**
* The metrics accumulator hands final totals to the program that launched the pipeline once the
* pipeline has ended. {@link #close()} also runs when the task is cancelled, which in streaming
* mode happens on every failover without the pipeline ending. Copying the metrics there is
* unnecessary and makes every JobManager status request merge and render the full metric state of
* each cancelled subtask. In streaming mode, only copy them once the input has ended.
*/
private boolean shouldRegisterMetricsForPipelineResult() {
return finished || !isStreaming;
}

void flushData() throws Exception {
// This is our last change to block shutdown of this operator while
// there are still remaining processing-time timers. Flink will ignore pending
Expand Down Expand Up @@ -689,6 +704,7 @@ void flushData() throws Exception {
public void finish() throws Exception {
try {
flushData();
finished = true;
} finally {
super.finish();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -240,6 +240,9 @@ public class DoFnOperator<PreInputT, InputT, OutputT>
/** Metrics container for reporting Beam metrics to Flink (null if metrics are disabled). */
transient @Nullable FlinkMetricContainer flinkMetricContainer;

/** Whether {@link #finish()} flushed all data, which only happens once the input has ended. */
private transient boolean finished;

/** Helper class to report the checkpoint duration. */
private transient @Nullable CheckpointStats checkpointStats;

Expand Down Expand Up @@ -299,7 +302,8 @@ public DoFnOperator(
this.sideInputTagMapping = sideInputTagMapping;
this.sideInputs = sideInputs;
this.serializedOptions = new SerializablePipelineOptions(options);
this.isStreaming = serializedOptions.get().as(FlinkPipelineOptions.class).isStreaming();
FlinkPipelineOptions flinkOptions = options.as(FlinkPipelineOptions.class);
this.isStreaming = flinkOptions.isStreaming();
this.windowingStrategy = windowingStrategy;
this.outputManagerFactory = outputManagerFactory;

Expand All @@ -311,8 +315,6 @@ public DoFnOperator(
this.timerCoder =
TimerInternals.TimerDataCoderV2.of(windowingStrategy.getWindowFn().windowCoder());

FlinkPipelineOptions flinkOptions = options.as(FlinkPipelineOptions.class);

this.maxBundleSize = flinkOptions.getMaxBundleSize();
Preconditions.checkArgument(maxBundleSize > 0, "Bundle size must be at least 1");
this.maxBundleTimeMills = flinkOptions.getMaxBundleTimeMills();
Expand Down Expand Up @@ -632,13 +634,26 @@ private void earlyBindStateIfNeeded() throws IllegalArgumentException, IllegalAc
}

void cleanUp() throws Exception {
Optional.ofNullable(flinkMetricContainer)
.ifPresent(FlinkMetricContainer::registerMetricsForPipelineResult);
if (shouldRegisterMetricsForPipelineResult()) {
Optional.ofNullable(flinkMetricContainer)
.ifPresent(FlinkMetricContainer::registerMetricsForPipelineResult);
}
Optional.ofNullable(checkFinishBundleTimer).ifPresent(timer -> timer.cancel(true));
Workarounds.deleteStaticCaches();
Optional.ofNullable(doFnInvoker).ifPresent(DoFnInvoker::invokeTeardown);
}

/**
* The metrics accumulator hands final totals to the program that launched the pipeline once the
* pipeline has ended. {@link #close()} also runs when the task is cancelled, which in streaming
* mode happens on every failover without the pipeline ending. Copying the metrics there is
* unnecessary and makes every JobManager status request merge and render the full metric state of
* each cancelled subtask. In streaming mode, only copy them once the input has ended.
*/
private boolean shouldRegisterMetricsForPipelineResult() {
return finished || !isStreaming;
}

void flushData() throws Exception {
// This is our last change to block shutdown of this operator while
// there are still remaining processing-time timers. Flink will ignore pending
Expand Down Expand Up @@ -691,6 +706,7 @@ void flushData() throws Exception {
public void finish() throws Exception {
try {
flushData();
finished = true;
} finally {
super.finish();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2343,12 +2343,7 @@ public void testAccumulatorRegistrationOnOperatorClose() throws Exception {

testHarness.open();

String metricContainerFieldName = "flinkMetricContainer";
FlinkMetricContainer monitoredContainer =
Mockito.spy(
(FlinkMetricContainer)
Whitebox.getInternalState(doFnOperator, metricContainerFieldName));
Whitebox.setInternalState(doFnOperator, metricContainerFieldName, monitoredContainer);
FlinkMetricContainer monitoredContainer = spyOnMetricContainer(doFnOperator);

// Closes and disposes the operator
testHarness.close();
Expand All @@ -2357,6 +2352,65 @@ public void testAccumulatorRegistrationOnOperatorClose() throws Exception {
Mockito.verify(monitoredContainer, Mockito.times(2)).registerMetricsForPipelineResult();
}

@Test
public void testAccumulatorRegistrationOnCancelInBatch() throws Exception {
DoFnOperator<String, String, String> doFnOperator = getOperatorForCleanupInspection();
OneInputStreamOperatorTestHarness<WindowedValue<String>, WindowedValue<String>> testHarness =
new OneInputStreamOperatorTestHarness<>(doFnOperator);

testHarness.open();

FlinkMetricContainer monitoredContainer = spyOnMetricContainer(doFnOperator);

// Closes the operator without finishing it, as when the task is cancelled.
doFnOperator.close();
Mockito.verify(monitoredContainer).registerMetricsForPipelineResult();
}

@Test
public void testAccumulatorRegistrationAfterFinishInStreaming() throws Exception {
FlinkPipelineOptions options = FlinkPipelineOptions.defaults();
options.setStreaming(true);
DoFnOperator<String, String, String> doFnOperator = getOperatorForCleanupInspection(options);
OneInputStreamOperatorTestHarness<WindowedValue<String>, WindowedValue<String>> testHarness =
new OneInputStreamOperatorTestHarness<>(doFnOperator);

testHarness.open();

FlinkMetricContainer monitoredContainer = spyOnMetricContainer(doFnOperator);

// Finishes and closes the operator, as when the input has ended.
testHarness.close();
Mockito.verify(monitoredContainer).registerMetricsForPipelineResult();
}

@Test
public void testNoAccumulatorRegistrationOnCancelInStreaming() throws Exception {
FlinkPipelineOptions options = FlinkPipelineOptions.defaults();
options.setStreaming(true);
DoFnOperator<String, String, String> doFnOperator = getOperatorForCleanupInspection(options);
OneInputStreamOperatorTestHarness<WindowedValue<String>, WindowedValue<String>> testHarness =
new OneInputStreamOperatorTestHarness<>(doFnOperator);

testHarness.open();

FlinkMetricContainer monitoredContainer = spyOnMetricContainer(doFnOperator);

// Closes the operator without finishing it, as when the task is cancelled.
doFnOperator.close();
Mockito.verify(monitoredContainer, Mockito.never()).registerMetricsForPipelineResult();
}

private static FlinkMetricContainer spyOnMetricContainer(DoFnOperator<?, ?, ?> doFnOperator) {
String metricContainerFieldName = "flinkMetricContainer";
FlinkMetricContainer monitoredContainer =
Mockito.spy(
(FlinkMetricContainer)
Whitebox.getInternalState(doFnOperator, metricContainerFieldName));
Whitebox.setInternalState(doFnOperator, metricContainerFieldName, monitoredContainer);
return monitoredContainer;
}

/**
* Ensures Jackson cache is cleaned to get rid of any references to the Flink Classloader. See
* https://jira.apache.org/jira/browse/BEAM-6460
Expand All @@ -2374,7 +2428,11 @@ public void testRemoveCachedClassReferences() throws Exception {
}

private static DoFnOperator<String, String, String> getOperatorForCleanupInspection() {
FlinkPipelineOptions options = FlinkPipelineOptions.defaults();
return getOperatorForCleanupInspection(FlinkPipelineOptions.defaults());
}

private static DoFnOperator<String, String, String> getOperatorForCleanupInspection(
FlinkPipelineOptions options) {
options.setParallelism(4);

TupleTag<String> outputTag = new TupleTag<>("main-output");
Expand Down
Loading