diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/WriteFiles.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/WriteFiles.java index b0b5051f3210..ee0f0bb20183 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/WriteFiles.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/WriteFiles.java @@ -29,6 +29,7 @@ import java.util.Map; import java.util.UUID; import java.util.concurrent.CompletionStage; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ThreadLocalRandom; import org.apache.beam.sdk.annotations.Internal; import org.apache.beam.sdk.coders.CannotProvideCoderException; @@ -46,6 +47,7 @@ import org.apache.beam.sdk.io.FileBasedSink.WriteOperation; import org.apache.beam.sdk.io.FileBasedSink.Writer; import org.apache.beam.sdk.io.fs.ResourceId; +import org.apache.beam.sdk.options.ExecutorOptions; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.ValueProvider; import org.apache.beam.sdk.options.ValueProvider.StaticValueProvider; @@ -1199,6 +1201,7 @@ public WriteShardsIntoTempFilesFn(Coder inputCoder) { private transient List> closeFutures = new ArrayList<>(); private transient List>> deferredOutput = new ArrayList<>(); + private transient ScheduledExecutorService executorService; // Ensure that transient fields are initialized. private void readObject(java.io.ObjectInputStream in) @@ -1208,6 +1211,11 @@ private void readObject(java.io.ObjectInputStream in) deferredOutput = new ArrayList<>(); } + @Setup + public void setup(PipelineOptions options) { + executorService = options.as(ExecutorOptions.class).getScheduledExecutorService(); + } + @ProcessElement public void processElement( ProcessContext c, BoundedWindow window, MultiOutputReceiver outputReceiver) @@ -1285,7 +1293,8 @@ private void closeWriterInBackground(Writer writer) { writer.cleanup(); throw e; } - })); + }, + executorService)); } @FinishBundle