Skip to content

Commit 8ec733e

Browse files
committed
Fix TSAN bug by swapping ForkJoinPool with ScheduledExecutorService
1 parent 58bac32 commit 8ec733e

2 files changed

Lines changed: 10 additions & 1 deletion

File tree

sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileBasedSink.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@
5555
import org.apache.beam.sdk.io.fs.MoveOptions.StandardMoveOptions;
5656
import org.apache.beam.sdk.io.fs.ResolveOptions.StandardResolveOptions;
5757
import org.apache.beam.sdk.io.fs.ResourceId;
58+
import org.apache.beam.sdk.options.ExecutorOptions;
5859
import org.apache.beam.sdk.options.PipelineOptions;
5960
import org.apache.beam.sdk.options.ValueProvider;
6061
import org.apache.beam.sdk.options.ValueProvider.NestedValueProvider;

sdks/java/core/src/main/java/org/apache/beam/sdk/io/WriteFiles.java

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
import java.util.Map;
3030
import java.util.UUID;
3131
import java.util.concurrent.CompletionStage;
32+
import java.util.concurrent.ScheduledExecutorService;
3233
import java.util.concurrent.ThreadLocalRandom;
3334
import org.apache.beam.sdk.annotations.Internal;
3435
import org.apache.beam.sdk.coders.CannotProvideCoderException;
@@ -1199,6 +1200,7 @@ public WriteShardsIntoTempFilesFn(Coder<UserT> inputCoder) {
11991200
private transient List<CompletionStage<Void>> closeFutures = new ArrayList<>();
12001201
private transient List<KV<Instant, FileResult<DestinationT>>> deferredOutput =
12011202
new ArrayList<>();
1203+
private transient ScheduledExecutorService executorService;
12021204

12031205
// Ensure that transient fields are initialized.
12041206
private void readObject(java.io.ObjectInputStream in)
@@ -1208,6 +1210,11 @@ private void readObject(java.io.ObjectInputStream in)
12081210
deferredOutput = new ArrayList<>();
12091211
}
12101212

1213+
@Setup
1214+
public void setup(PipelineOptions options) {
1215+
executorService = options.as(ExecutorOptions.class).getScheduledExecutorService();
1216+
}
1217+
12111218
@ProcessElement
12121219
public void processElement(
12131220
ProcessContext c, BoundedWindow window, MultiOutputReceiver outputReceiver)
@@ -1285,7 +1292,8 @@ private void closeWriterInBackground(Writer<DestinationT, OutputT> writer) {
12851292
writer.cleanup();
12861293
throw e;
12871294
}
1288-
}));
1295+
},
1296+
executorService));
12891297
}
12901298

12911299
@FinishBundle

0 commit comments

Comments
 (0)