This is an automated email from the ASF dual-hosted git repository.
reuvenlax pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 086ab9bc601 Merge pull request #40324 from
reuvenlax/improve_test_presubmit
086ab9bc601 is described below
commit 086ab9bc6013319e104d622517988602bafa842f
Author: Reuven Lax <[email protected]>
AuthorDate: Thu Oct 1 12:30:11 2026 -0700
Merge pull request #40324 from reuvenlax/improve_test_presubmit
Make Java PreSubmit faster: Parallelize test
---
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 1 +
.../apache/beam/examples/WindowedWordCountIT.java | 4 +--
runners/direct-java/build.gradle | 3 ++
.../runners/dataflow/TestDataflowRunnerTest.java | 4 +--
sdks/java/core/build.gradle | 1 +
.../beam/sdk/testing/BeamParallelJunit4Runner.java | 1 +
.../java/org/apache/beam/sdk/testing/PAssert.java | 15 +++++----
.../java/org/apache/beam/sdk/io/FileIOTest.java | 4 +--
.../apache/beam/sdk/io/GenerateSequenceTest.java | 4 +--
.../org/apache/beam/sdk/io/TFRecordIOTest.java | 4 +--
.../io/TFRecordSchemaTransformProviderTest.java | 4 +--
.../org/apache/beam/sdk/io/TextIOReadTest.java | 9 +++--
.../org/apache/beam/sdk/io/WriteFilesTest.java | 5 +--
.../beam/sdk/schemas/AutoValueSchemaTest.java | 4 +--
.../beam/sdk/schemas/transforms/ConvertTest.java | 4 +--
.../sdk/transforms/ApproximateQuantilesTest.java | 4 +--
.../org/apache/beam/sdk/transforms/WaitTest.java | 25 +++++++++-----
.../beam/sdk/values/PCollectionViewsTest.java | 4 +--
.../python/PythonExternalTransformTest.java | 1 +
.../sorter/BufferedExternalSorterTest.java | 4 +--
.../beam/fn/harness/status/MemoryMonitor.java | 38 +++++++++++++++++-----
.../beam/fn/harness/status/MemoryMonitorTest.java | 11 +++++++
.../io/contextualtextio/ContextualTextIOTest.java | 9 +++--
.../sdk/io/sparkreceiver/SparkReceiverIOTest.java | 17 ++++++----
.../io/synthetic/SyntheticBoundedSourceTest.java | 4 +--
.../beam/sdk/io/synthetic/SyntheticStepTest.java | 4 +--
sdks/java/testing/nexmark/build.gradle | 5 +++
.../nexmark/queries/BoundedSideInputJoinTest.java | 4 +--
.../apache/beam/sdk/nexmark/queries/QueryTest.java | 4 +--
.../nexmark/queries/SessionSideInputJoinTest.java | 4 +--
.../beam/sdk/nexmark/queries/SqlQueryTest.java | 3 +-
.../queries/sql/SqlBoundedSideInputJoinTest.java | 3 +-
32 files changed, 136 insertions(+), 75 deletions(-)
diff --git
a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
index 0c34e1529e8..c66213fe750 100644
--- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
+++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
@@ -1259,6 +1259,7 @@ class BeamModulePlugin implements Plugin<Project> {
useJUnit {}
// default maxHeapSize on gradle 5 is 512m, lets increase to handle
more demanding tests
maxHeapSize = '2g'
+ systemProperty 'beam.test.parallelThreads',
project.findProperty('testParallelThreads') ?: '8'
// Windows OS: Snappy needs an executable temp dir for native lib.
Default AppData/Temp
// failing with Access error without elevated permissions
if (System.getProperty("os.name").toLowerCase().contains("windows")) {
diff --git
a/examples/java/src/test/java/org/apache/beam/examples/WindowedWordCountIT.java
b/examples/java/src/test/java/org/apache/beam/examples/WindowedWordCountIT.java
index f7b858ccf5e..8bc0a3a8ea0 100644
---
a/examples/java/src/test/java/org/apache/beam/examples/WindowedWordCountIT.java
+++
b/examples/java/src/test/java/org/apache/beam/examples/WindowedWordCountIT.java
@@ -36,6 +36,7 @@ import
org.apache.beam.sdk.io.fs.ResolveOptions.StandardResolveOptions;
import org.apache.beam.sdk.io.fs.ResourceId;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.options.StreamingOptions;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.SerializableMatcher;
import org.apache.beam.sdk.testing.StreamingIT;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -59,10 +60,9 @@ import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.rules.TestName;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** End-to-end integration test of {@link WindowedWordCount}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class WindowedWordCountIT {
@Rule public TestName testName = new TestName();
diff --git a/runners/direct-java/build.gradle b/runners/direct-java/build.gradle
index 1ab702da321..547f49ce910 100644
--- a/runners/direct-java/build.gradle
+++ b/runners/direct-java/build.gradle
@@ -120,6 +120,7 @@ task needsRunnerTests(type: Test) {
group = "Verification"
description = "Runs tests that require a runner to validate that
pipelines/transforms work correctly"
+ maxParallelForks 4
testLogging.showStandardStreams = true
String[] pipelineOptions = ["--runner=DirectRunner",
"--runnerDeterminedSharding=false"]
@@ -160,6 +161,8 @@ task validatesRunner(type: Test) {
group = "Verification"
description "Validates Direct runner"
+ maxParallelForks 4
+
String[] pipelineOptions = ["--runner=DirectRunner",
"--runnerDeterminedSharding=false"]
systemProperty "beamTestPipelineOptions",
pipelineOptionsStringCrossPlatformHandling(pipelineOptions)
diff --git
a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/TestDataflowRunnerTest.java
b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/TestDataflowRunnerTest.java
index 4f6ff01c327..5cd16757e36 100644
---
a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/TestDataflowRunnerTest.java
+++
b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/TestDataflowRunnerTest.java
@@ -51,6 +51,7 @@ import
org.apache.beam.sdk.extensions.gcp.storage.NoopPathValidator;
import org.apache.beam.sdk.extensions.gcp.util.Transport;
import org.apache.beam.sdk.io.FileSystems;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.SerializableMatcher;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -69,13 +70,12 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.ExpectedException;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
import org.mockito.Mock;
import org.mockito.Mockito;
import org.mockito.MockitoAnnotations;
/** Tests for {@link TestDataflowRunner}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class TestDataflowRunnerTest {
@Rule public ExpectedException expectedException = ExpectedException.none();
@Mock private DataflowClient mockClient;
diff --git a/sdks/java/core/build.gradle b/sdks/java/core/build.gradle
index f95a8a53906..2bb8f4bcd0b 100644
--- a/sdks/java/core/build.gradle
+++ b/sdks/java/core/build.gradle
@@ -67,6 +67,7 @@ processResources {
// Exclude tests that need a runner
test {
+ maxParallelForks 4
systemProperty "beamUseDummyRunner", "true"
useJUnit {
excludeCategories "org.apache.beam.sdk.testing.NeedsRunner"
diff --git
a/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/BeamParallelJunit4Runner.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/BeamParallelJunit4Runner.java
index bcaead005ad..59527b161f0 100644
---
a/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/BeamParallelJunit4Runner.java
+++
b/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/BeamParallelJunit4Runner.java
@@ -182,6 +182,7 @@ public final class BeamParallelJunit4Runner extends
BlockJUnit4ClassRunner {
if (clazz == null || clazz == Object.class) {
return false;
}
+ // Don't access CLASS_SERIAL_CACHE here, as ConcurrentHashMap does not
support reentrancy.
return clazz.isAnnotationPresent(SerialTest.class)
|| computeClassSerial(clazz.getSuperclass())
|| computeClassSerial(clazz.getEnclosingClass());
diff --git
a/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/PAssert.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/PAssert.java
index be834228fbb..a4cbb29103c 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/PAssert.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/PAssert.java
@@ -32,6 +32,7 @@ import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.NoSuchElementException;
+import java.util.concurrent.atomic.AtomicInteger;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.Pipeline.PipelineVisitor;
import org.apache.beam.sdk.PipelineRunner;
@@ -121,10 +122,12 @@ public class PAssert {
private static final Counter failureCounter =
Metrics.counter(PAssert.class, PAssert.FAILURE_COUNTER);
- private static int assertCount = 0;
+ // Atomic so that PAssert transforms constructed concurrently (e.g. by tests
running under
+ // BeamParallelJunit4Runner) never receive duplicate names within a pipeline.
+ private static final AtomicInteger assertCount = new AtomicInteger(0);
private static String nextAssertionName() {
- return "PAssert$" + assertCount++;
+ return "PAssert$" + assertCount.getAndIncrement();
}
// Do not instantiate.
@@ -786,7 +789,7 @@ public class PAssert {
SerializableFunction<Iterable<T>, Void> checkerFn =
(SerializableFunction) new MatcherCheckerFn<>(matcher);
actual.apply(
- "PAssert$" + assertCount++,
+ nextAssertionName(),
new GroupThenAssert<>(checkerFn, rewindowingStrategy, paneExtractor,
site));
return this;
}
@@ -940,7 +943,7 @@ public class PAssert {
public PCollectionSingletonIterableAssert<T> satisfies(
SerializableFunction<Iterable<T>, Void> checkerFn) {
actual.apply(
- "PAssert$" + assertCount++,
+ nextAssertionName(),
new GroupThenAssertForSingleton<>(checkerFn, rewindowingStrategy,
paneExtractor, site));
return this;
}
@@ -1034,7 +1037,7 @@ public class PAssert {
@Override
public PCollectionSingletonAssert<T> satisfies(SerializableFunction<T,
Void> checkerFn) {
actual.apply(
- "PAssert$" + assertCount++,
+ nextAssertionName(),
new GroupThenAssertForSingleton<>(checkerFn, rewindowingStrategy,
paneExtractor, site));
return this;
}
@@ -1173,7 +1176,7 @@ public class PAssert {
actual
.getPipeline()
.apply(
- "PAssert$" + assertCount++,
+ nextAssertionName(),
new OneSideInputAssert<>(
CreateActual.from(actual, rewindowActuals, paneExtractor,
view),
rewindowActuals.windowDummy(),
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java
index c5a227d46b8..7a0dfb4822b 100644
--- a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java
+++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java
@@ -56,6 +56,7 @@ import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.state.StateSpec;
import org.apache.beam.sdk.state.StateSpecs;
import org.apache.beam.sdk.state.ValueState;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -91,10 +92,9 @@ import org.junit.rules.ExpectedException;
import org.junit.rules.TemporaryFolder;
import org.junit.rules.Timeout;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Tests for {@link FileIO}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class FileIOTest implements Serializable {
@Rule public transient TestPipeline p = TestPipeline.create();
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/GenerateSequenceTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/GenerateSequenceTest.java
index 2cb15797758..25e48563410 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/GenerateSequenceTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/GenerateSequenceTest.java
@@ -21,6 +21,7 @@ import static
org.apache.beam.sdk.transforms.display.DisplayDataMatchers.hasDisp
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.is;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -40,10 +41,9 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Tests for {@link GenerateSequence}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class GenerateSequenceTest {
public static void addCountingAsserts(PCollection<Long> input, long start,
long end) {
// Count == numElements
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TFRecordIOTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TFRecordIOTest.java
index a38faf077e0..6cdb1c49d5e 100644
--- a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TFRecordIOTest.java
+++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TFRecordIOTest.java
@@ -59,6 +59,7 @@ import org.apache.beam.sdk.coders.StringUtf8Coder;
import org.apache.beam.sdk.io.FileIO.ReadableFile;
import org.apache.beam.sdk.io.TFRecordIO.TFRecordCodec;
import org.apache.beam.sdk.io.fs.MatchResult.Metadata;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -77,10 +78,9 @@ import org.junit.experimental.categories.Category;
import org.junit.rules.ExpectedException;
import org.junit.rules.TemporaryFolder;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Tests for TFRecordIO Read and Write transforms. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class TFRecordIOTest {
/*
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TFRecordSchemaTransformProviderTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TFRecordSchemaTransformProviderTest.java
index 65e38ed9d29..7d77c477a53 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TFRecordSchemaTransformProviderTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TFRecordSchemaTransformProviderTest.java
@@ -50,6 +50,7 @@ import org.apache.beam.sdk.schemas.Schema;
import org.apache.beam.sdk.schemas.transforms.SchemaTransform;
import org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider;
import org.apache.beam.sdk.schemas.transforms.providers.ErrorHandling;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -67,10 +68,9 @@ import org.junit.experimental.categories.Category;
import org.junit.rules.ExpectedException;
import org.junit.rules.TemporaryFolder;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Tests for TFRecordIO Read and Write transforms. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class TFRecordSchemaTransformProviderTest {
/*
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TextIOReadTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TextIOReadTest.java
index e7d25608eb7..ce7b62fe1d2 100644
--- a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TextIOReadTest.java
+++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/TextIOReadTest.java
@@ -73,6 +73,7 @@ import org.apache.beam.sdk.io.fs.MatchResult.Metadata;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.options.ValueProvider;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.SourceTestUtils;
@@ -106,7 +107,6 @@ import org.junit.experimental.categories.Category;
import org.junit.experimental.runners.Enclosed;
import org.junit.rules.TemporaryFolder;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
import org.junit.runners.Parameterized;
/** Tests for {@link TextIO.Read}. */
@@ -633,7 +633,7 @@ public class TextIOReadTest {
}
/** Tests for some basic operations in {@link TextIO.Read}. */
- @RunWith(JUnit4.class)
+ @RunWith(BeamParallelJunit4Runner.class)
public static class BasicIOTest {
@Rule public TemporaryFolder tempFolder = new TemporaryFolder();
@Rule public TestPipeline p = TestPipeline.create();
@@ -1171,6 +1171,9 @@ public class TextIOReadTest {
@Test
@Category({NeedsRunner.class, UsesUnboundedSplittableParDo.class})
+ // The watch terminates after 3s without new output; running concurrently
with other pipelines
+ // in this JVM can stall the in-pipeline writer long enough to truncate
the results.
+ @BeamParallelJunit4Runner.SerialTest
public void testReadWatchForNewFiles() throws IOException,
InterruptedException {
final Path basePath = tempFolder.getRoot().toPath().resolve("readWatch");
basePath.toFile().mkdir();
@@ -1202,7 +1205,7 @@ public class TextIOReadTest {
}
/** Tests for TextSource class. */
- @RunWith(JUnit4.class)
+ @RunWith(BeamParallelJunit4Runner.class)
public static class TextSourceTest {
@Rule public transient TestPipeline pipeline = TestPipeline.create();
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/WriteFilesTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/WriteFilesTest.java
index 78a9120517e..09f3f04d5b0 100644
--- a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/WriteFilesTest.java
+++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/WriteFilesTest.java
@@ -67,6 +67,7 @@ import org.apache.beam.sdk.io.fs.ResourceId;
import org.apache.beam.sdk.options.Description;
import
org.apache.beam.sdk.options.PipelineOptionsFactoryTest.TestPipelineOptions;
import org.apache.beam.sdk.options.ValueProvider.StaticValueProvider;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -116,10 +117,9 @@ import org.junit.experimental.categories.Category;
import org.junit.rules.ExpectedException;
import org.junit.rules.TemporaryFolder;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Tests for the WriteFiles PTransform. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class WriteFilesTest {
@Rule public TemporaryFolder tmpFolder = new TemporaryFolder();
@Rule public final TestPipeline p = TestPipeline.create();
@@ -633,6 +633,7 @@ public class WriteFilesTest {
@Test
@Category(NeedsRunner.class)
+ @BeamParallelJunit4Runner.SerialTest
public void testWriteEvictWritersWhenFullCloseExceptionCleanup() {
FailingCloseOnEvictSink.EVICTED_TEMP_FILE.set(null);
FailingCloseOnEvictSink.THREW_ON_CLOSE.set(false);
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/AutoValueSchemaTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/AutoValueSchemaTest.java
index 99b858a4e30..00e2d7628e3 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/AutoValueSchemaTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/AutoValueSchemaTest.java
@@ -40,6 +40,7 @@ import
org.apache.beam.sdk.schemas.annotations.SchemaFieldDescription;
import org.apache.beam.sdk.schemas.annotations.SchemaFieldName;
import org.apache.beam.sdk.schemas.annotations.SchemaFieldNumber;
import org.apache.beam.sdk.schemas.utils.SchemaTestUtils;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.transforms.SerializableFunction;
import org.apache.beam.sdk.util.SerializableUtils;
import org.apache.beam.sdk.values.Row;
@@ -51,10 +52,9 @@ import org.joda.time.DateTime;
import org.joda.time.Instant;
import org.junit.Test;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Tests for {@link AutoValueSchema}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class AutoValueSchemaTest {
static final DateTime DATE = DateTime.parse("1979-03-14");
static final byte[] BYTE_ARRAY =
"bytearray".getBytes(StandardCharsets.UTF_8);
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/transforms/ConvertTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/transforms/ConvertTest.java
index c206bd8f61f..8a20c6d5fac 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/transforms/ConvertTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/transforms/ConvertTest.java
@@ -26,6 +26,7 @@ import org.apache.beam.sdk.schemas.Schema;
import org.apache.beam.sdk.schemas.Schema.FieldType;
import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
import org.apache.beam.sdk.schemas.logicaltypes.SqlTypes;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -42,10 +43,9 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Tests for the {@link Convert} class. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
@Category(UsesSchema.class)
public class ConvertTest {
@Rule public final transient TestPipeline pipeline = TestPipeline.create();
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ApproximateQuantilesTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ApproximateQuantilesTest.java
index d5949558f27..37da2bcdd5c 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ApproximateQuantilesTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ApproximateQuantilesTest.java
@@ -33,6 +33,7 @@ import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.coders.BigEndianIntegerCoder;
import org.apache.beam.sdk.coders.KvCoder;
import org.apache.beam.sdk.coders.StringUtf8Coder;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -50,14 +51,13 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
import org.junit.runners.Parameterized;
/** Tests for {@link ApproximateQuantiles}. */
public class ApproximateQuantilesTest {
/** Tests for the overall combiner behavior. */
- @RunWith(JUnit4.class)
+ @RunWith(BeamParallelJunit4Runner.class)
public static class CombinerTests {
static final List<KV<String, Integer>> TABLE =
Arrays.asList(
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/WaitTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/WaitTest.java
index 4c1c692765f..4dea35294bd 100644
--- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/WaitTest.java
+++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/WaitTest.java
@@ -33,6 +33,7 @@ import org.apache.beam.sdk.schemas.NoSuchSchemaException;
import org.apache.beam.sdk.schemas.SchemaCoder;
import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
import org.apache.beam.sdk.schemas.annotations.SchemaCreate;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -64,10 +65,9 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Tests for {@link Wait}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class WaitTest implements Serializable {
@Rule public transient TestPipeline p = TestPipeline.create();
@@ -167,8 +167,8 @@ public class WaitTest implements Serializable {
return p.apply(name, stream.advanceWatermarkToInfinity());
}
- private static final AtomicReference<Instant> TEST_WAIT_MAX_MAIN_TIMESTAMP =
- new AtomicReference<>();
+ private static final java.util.Map<java.util.UUID, AtomicReference<Instant>>
+ TEST_WAIT_MAX_MAIN_TIMESTAMP = new
java.util.concurrent.ConcurrentHashMap<>();
@Test
@Category({NeedsRunner.class, UsesTestStreamWithProcessingTime.class})
@@ -273,6 +273,7 @@ public class WaitTest implements Serializable {
}
@Test
+ @BeamParallelJunit4Runner.SerialTest
@Category({NeedsRunner.class, UsesTestStream.class})
public void testWindowExpiration() throws NoSuchSchemaException {
PROCESSED_LONGS.clear();
@@ -399,7 +400,8 @@ public class WaitTest implements Serializable {
@Nullable WindowFn<? super Long, ?> mainWindowFn,
int numSignalElements,
@Nullable WindowFn<? super Long, ?> signalWindowFn) {
- TEST_WAIT_MAX_MAIN_TIMESTAMP.set(null);
+ final java.util.UUID testId = java.util.UUID.randomUUID();
+ TEST_WAIT_MAX_MAIN_TIMESTAMP.put(testId, new AtomicReference<>());
Instant base = Instant.now();
@@ -442,7 +444,10 @@ public class WaitTest implements Serializable {
new DoFn<Long, Long>() {
@ProcessElement
public void process(ProcessContext c) {
- Instant maxMainTimestamp =
TEST_WAIT_MAX_MAIN_TIMESTAMP.get();
+ Instant maxMainTimestamp =
+ TEST_WAIT_MAX_MAIN_TIMESTAMP
+ .computeIfAbsent(testId, unused -> new
AtomicReference<>())
+ .get();
if (maxMainTimestamp != null) {
assertFalse(
"Signal at timestamp "
@@ -463,14 +468,16 @@ public class WaitTest implements Serializable {
new DoFn<Long, Long>() {
@ProcessElement
public void process(ProcessContext c,
@SuppressWarnings("unused") BoundedWindow w) {
+ AtomicReference<Instant> ref =
+ TEST_WAIT_MAX_MAIN_TIMESTAMP.computeIfAbsent(
+ testId, unused -> new AtomicReference<>());
while (true) {
- Instant maxMainTimestamp =
TEST_WAIT_MAX_MAIN_TIMESTAMP.get();
+ Instant maxMainTimestamp = ref.get();
Instant newMaxTimestamp =
(maxMainTimestamp == null ||
c.timestamp().isAfter(maxMainTimestamp))
? c.timestamp()
: maxMainTimestamp;
- if (TEST_WAIT_MAX_MAIN_TIMESTAMP.compareAndSet(
- maxMainTimestamp, newMaxTimestamp)) {
+ if (ref.compareAndSet(maxMainTimestamp, newMaxTimestamp)) {
break;
}
}
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/values/PCollectionViewsTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/values/PCollectionViewsTest.java
index 457c83f9bed..943cf354fd5 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/values/PCollectionViewsTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/values/PCollectionViewsTest.java
@@ -33,16 +33,16 @@ import java.util.Random;
import java.util.TreeSet;
import java.util.stream.IntStream;
import org.apache.beam.sdk.io.range.OffsetRange;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ArrayListMultimap;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ListMultimap;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
import org.junit.Test;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Tests for {@link PCollectionViews}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class PCollectionViewsTest {
@Test
public void testEmpty() {
diff --git
a/sdks/java/extensions/python/src/test/java/org/apache/beam/sdk/extensions/python/PythonExternalTransformTest.java
b/sdks/java/extensions/python/src/test/java/org/apache/beam/sdk/extensions/python/PythonExternalTransformTest.java
index 30fe0b90f39..fc3adb66afc 100644
---
a/sdks/java/extensions/python/src/test/java/org/apache/beam/sdk/extensions/python/PythonExternalTransformTest.java
+++
b/sdks/java/extensions/python/src/test/java/org/apache/beam/sdk/extensions/python/PythonExternalTransformTest.java
@@ -377,6 +377,7 @@ public class PythonExternalTransformTest implements
Serializable {
}
@Test
+ @Category({ValidatesRunner.class, UsesPythonExpansionService.class})
public void testLoopbackEnvironmentWithPythonExternalTransform() {
PortablePipelineOptions options =
PipelineOptionsFactory.create().as(PortablePipelineOptions.class);
diff --git
a/sdks/java/extensions/sorter/src/test/java/org/apache/beam/sdk/extensions/sorter/BufferedExternalSorterTest.java
b/sdks/java/extensions/sorter/src/test/java/org/apache/beam/sdk/extensions/sorter/BufferedExternalSorterTest.java
index 5d6e54c6b68..bce5e02a659 100644
---
a/sdks/java/extensions/sorter/src/test/java/org/apache/beam/sdk/extensions/sorter/BufferedExternalSorterTest.java
+++
b/sdks/java/extensions/sorter/src/test/java/org/apache/beam/sdk/extensions/sorter/BufferedExternalSorterTest.java
@@ -33,6 +33,7 @@ import java.nio.file.Path;
import java.nio.file.SimpleFileVisitor;
import java.nio.file.attribute.BasicFileAttributes;
import java.util.Arrays;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.values.KV;
import org.checkerframework.checker.nullness.qual.EnsuresNonNull;
import org.checkerframework.checker.nullness.qual.Nullable;
@@ -42,10 +43,9 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.ExpectedException;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Tests for {@link BufferedExternalSorter}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
@SuppressWarnings({
"rawtypes" // TODO(https://github.com/apache/beam/issues/20447)
})
diff --git
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/status/MemoryMonitor.java
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/status/MemoryMonitor.java
index 9dbd1c96359..797ec1866aa 100644
---
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/status/MemoryMonitor.java
+++
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/status/MemoryMonitor.java
@@ -200,6 +200,19 @@ public class MemoryMonitor implements Runnable {
@VisibleForTesting final boolean gzipCompress;
+ /** Interface for dumping the heap to a file (for testing). */
+ @VisibleForTesting
+ interface HeapDumper {
+ void dump(File destination)
+ throws MalformedObjectNameException,
+ InstanceNotFoundException,
+ ReflectionException,
+ MBeanException,
+ IOException;
+ }
+
+ @VisibleForTesting HeapDumper heapDumper = MemoryMonitor::dumpJvmHeap;
+
public static MemoryMonitor fromOptions(PipelineOptions options) {
SdkHarnessOptions sdkHarnessOptions = options.as(SdkHarnessOptions.class);
@Nullable String uploadFilePath =
sdkHarnessOptions.getRemoteHeapDumpLocation();
@@ -627,7 +640,20 @@ public class MemoryMonitor implements Runnable {
IOException {
Preconditions.checkState(
canDumpHeap, "Bug! Attempt to dump heap even though it should be
disabled.");
- return dumpHeap(localDumpFolder);
+ return dumpHeap(localDumpFolder, heapDumper);
+ }
+
+ private static void dumpJvmHeap(File fileName)
+ throws MalformedObjectNameException,
+ InstanceNotFoundException,
+ ReflectionException,
+ MBeanException {
+ boolean liveObjectsOnly = false;
+ MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
+ ObjectName oname = new
ObjectName("com.sun.management:type=HotSpotDiagnostic");
+ Object[] parameters = {fileName.getPath(), liveObjectsOnly};
+ String[] signatures = {String.class.getName(), boolean.class.getName()};
+ mbs.invoke(oname, "dumpHeap", parameters, signatures);
}
/**
@@ -636,24 +662,18 @@ public class MemoryMonitor implements Runnable {
* <p>NOTE: We deliberately don't salt the heap dump filename so as to
minimize disk impact of
* repeated dumps. These files can be of comparable size to the local disk.
*/
- private static synchronized File dumpHeap(File directory)
+ private static synchronized File dumpHeap(File directory, HeapDumper
heapDumper)
throws MalformedObjectNameException,
InstanceNotFoundException,
ReflectionException,
MBeanException,
IOException {
-
- boolean liveObjectsOnly = false;
File fileName = new File(directory, "heap_dump.hprof");
if (fileName.exists() && !fileName.delete()) {
throw new IOException("heap_dump.hprof already existed and couldn't be
deleted!");
}
- MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
- ObjectName oname = new
ObjectName("com.sun.management:type=HotSpotDiagnostic");
- Object[] parameters = {fileName.getPath(), liveObjectsOnly};
- String[] signatures = {String.class.getName(), boolean.class.getName()};
- mbs.invoke(oname, "dumpHeap", parameters, signatures);
+ heapDumper.dump(fileName);
if
(java.nio.file.FileSystems.getDefault().supportedFileAttributeViews().contains("posix"))
{
Files.setPosixFilePermissions(
diff --git
a/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/status/MemoryMonitorTest.java
b/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/status/MemoryMonitorTest.java
index f28df4f2d10..199b9cee06d 100644
---
a/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/status/MemoryMonitorTest.java
+++
b/sdks/java/harness/src/test/java/org/apache/beam/fn/harness/status/MemoryMonitorTest.java
@@ -26,6 +26,8 @@ import static org.junit.Assert.assertTrue;
import java.io.File;
import java.io.FileInputStream;
import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
@@ -47,6 +49,10 @@ import org.junit.runners.JUnit4;
@RunWith(JUnit4.class)
public class MemoryMonitorTest {
+ private static final MemoryMonitor.HeapDumper FAKE_HEAP_DUMPER =
+ destination ->
+ Files.write(destination.toPath(), "fake heap
dump".getBytes(StandardCharsets.UTF_8));
+
@Rule public TemporaryFolder tempFolder = new TemporaryFolder();
static class FakeGCStatsProvider implements MemoryMonitor.GCStatsProvider {
@@ -120,6 +126,7 @@ public class MemoryMonitorTest {
public void heapDumpTwice() throws Exception {
MemoryMonitor monitor =
MemoryMonitor.forTest(provider, 10, 0, true, 50.0, null,
localDumpFolder, false);
+ monitor.heapDumper = FAKE_HEAP_DUMPER;
File dump1 = monitor.dumpHeap();
assertNotNull(dump1);
assertTrue(dump1.exists());
@@ -137,6 +144,7 @@ public class MemoryMonitorTest {
MemoryMonitor monitor =
MemoryMonitor.forTest(
provider, 10, 0, true, 50.0, remoteFolder.getPath(),
localDumpFolder, false);
+ monitor.heapDumper = FAKE_HEAP_DUMPER;
// Force the monitor to generate a local heap dump
monitor.dumpHeap();
@@ -156,6 +164,7 @@ public class MemoryMonitorTest {
MemoryMonitor monitor =
MemoryMonitor.forTest(
provider, 10, 0, true, 50.0, remoteFolder.getPath(),
localDumpFolder, true);
+ monitor.heapDumper = FAKE_HEAP_DUMPER;
// Force the monitor to generate a local heap dump
monitor.dumpHeap();
@@ -178,6 +187,7 @@ public class MemoryMonitorTest {
public void uploadFileDisabled() throws Exception {
MemoryMonitor monitor =
MemoryMonitor.forTest(provider, 10, 0, true, 50.0, null,
localDumpFolder, false);
+ monitor.heapDumper = FAKE_HEAP_DUMPER;
// Force the monitor to generate a local heap dump
monitor.dumpHeap();
@@ -306,6 +316,7 @@ public class MemoryMonitorTest {
PipelineOptionsFactory.fromArgs(
"--enableHeapDumps", "--remoteHeapDumpLocation=" +
remoteFolder)
.create());
+ m.heapDumper = FAKE_HEAP_DUMPER;
assertTrue(m.canDumpHeap);
assertEquals(subfolder, m.localDumpFolder);
assertEquals(remoteFolder.toString(), m.uploadFilePath);
diff --git
a/sdks/java/io/contextualtextio/src/test/java/org/apache/beam/sdk/io/contextualtextio/ContextualTextIOTest.java
b/sdks/java/io/contextualtextio/src/test/java/org/apache/beam/sdk/io/contextualtextio/ContextualTextIOTest.java
index bce7db2f264..82956316be7 100644
---
a/sdks/java/io/contextualtextio/src/test/java/org/apache/beam/sdk/io/contextualtextio/ContextualTextIOTest.java
+++
b/sdks/java/io/contextualtextio/src/test/java/org/apache/beam/sdk/io/contextualtextio/ContextualTextIOTest.java
@@ -69,6 +69,7 @@ import org.apache.beam.sdk.io.fs.ResourceId;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.options.ValueProvider;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.SourceTestUtils;
@@ -99,7 +100,6 @@ import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.rules.TemporaryFolder;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
import org.junit.runners.Parameterized;
/** Tests for {@link ContextualTextIO.Read}. */
@@ -507,7 +507,7 @@ public class ContextualTextIOTest {
}
/** Tests Specific for checking functionality of ContextualTextIO. */
- @RunWith(JUnit4.class)
+ @RunWith(BeamParallelJunit4Runner.class)
public static class ContextualTextIOSpecificTests {
@Rule public TemporaryFolder tempFolder = new TemporaryFolder();
@Rule public TestPipeline p = TestPipeline.create();
@@ -787,7 +787,7 @@ public class ContextualTextIOTest {
}
/** Tests for some basic operations in {@link ContextualTextIO.Read}. */
- @RunWith(JUnit4.class)
+ @RunWith(BeamParallelJunit4Runner.class)
public static class BasicIOTest {
@Rule public TemporaryFolder tempFolder = new TemporaryFolder();
@Rule public TestPipeline p = TestPipeline.create();
@@ -1276,6 +1276,9 @@ public class ContextualTextIOTest {
@Test
@Category({NeedsRunner.class, UsesUnboundedSplittableParDo.class})
+ // The watch terminates after 3s without new output; running concurrently
with other pipelines
+ // in this JVM can stall the in-pipeline writer long enough to truncate
the results.
+ @BeamParallelJunit4Runner.SerialTest
public void testReadWatchForNewFiles() throws IOException,
InterruptedException {
final Path basePath = tempFolder.getRoot().toPath().resolve("readWatch");
basePath.toFile().mkdir();
diff --git
a/sdks/java/io/sparkreceiver/3/src/test/java/org/apache/beam/sdk/io/sparkreceiver/SparkReceiverIOTest.java
b/sdks/java/io/sparkreceiver/3/src/test/java/org/apache/beam/sdk/io/sparkreceiver/SparkReceiverIOTest.java
index bb482e79838..c924543fe66 100644
---
a/sdks/java/io/sparkreceiver/3/src/test/java/org/apache/beam/sdk/io/sparkreceiver/SparkReceiverIOTest.java
+++
b/sdks/java/io/sparkreceiver/3/src/test/java/org/apache/beam/sdk/io/sparkreceiver/SparkReceiverIOTest.java
@@ -23,6 +23,7 @@ import static org.junit.Assert.assertThrows;
import java.util.ArrayList;
import java.util.List;
import org.apache.beam.sdk.coders.StringUtf8Coder;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
import org.apache.beam.sdk.testing.TestPipelineOptions;
@@ -33,23 +34,23 @@ import org.joda.time.Instant;
import org.junit.Rule;
import org.junit.Test;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Test class for {@link SparkReceiverIO}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class SparkReceiverIOTest {
- public static final TestPipelineOptions OPTIONS =
- TestPipeline.testingPipelineOptions().as(TestPipelineOptions.class);
public static final long PULL_FREQUENCY_SEC = 1L;
public static final long START_POLL_TIMEOUT_SEC = 2L;
public static final long START_OFFSET = 0L;
- static {
- OPTIONS.setBlockOnRun(false);
+ private static TestPipelineOptions createOptions() {
+ TestPipelineOptions options =
+ TestPipeline.testingPipelineOptions().as(TestPipelineOptions.class);
+ options.setBlockOnRun(false);
+ return options;
}
- @Rule public final transient TestPipeline pipeline =
TestPipeline.fromOptions(OPTIONS);
+ @Rule public final transient TestPipeline pipeline =
TestPipeline.fromOptions(createOptions());
@Test
public void testReadBuildsCorrectly() {
@@ -126,6 +127,7 @@ public class SparkReceiverIOTest {
}
@Test
+ @BeamParallelJunit4Runner.SerialTest
public void testReadFromCustomReceiverWithOffset() {
CustomReceiverWithOffset.shouldFailInTheMiddle = false;
ReceiverBuilder<String, CustomReceiverWithOffset> receiverBuilder =
@@ -150,6 +152,7 @@ public class SparkReceiverIOTest {
}
@Test
+ @BeamParallelJunit4Runner.SerialTest
public void testReadFromCustomReceiverWithOffsetFailsAndReread() {
CustomReceiverWithOffset.shouldFailInTheMiddle = true;
ReceiverBuilder<String, CustomReceiverWithOffset> receiverBuilder =
diff --git
a/sdks/java/io/synthetic/src/test/java/org/apache/beam/sdk/io/synthetic/SyntheticBoundedSourceTest.java
b/sdks/java/io/synthetic/src/test/java/org/apache/beam/sdk/io/synthetic/SyntheticBoundedSourceTest.java
index c6c8df4d8a5..3cb63c3ee3b 100644
---
a/sdks/java/io/synthetic/src/test/java/org/apache/beam/sdk/io/synthetic/SyntheticBoundedSourceTest.java
+++
b/sdks/java/io/synthetic/src/test/java/org/apache/beam/sdk/io/synthetic/SyntheticBoundedSourceTest.java
@@ -30,6 +30,7 @@ import org.apache.beam.sdk.io.BoundedSource;
import org.apache.beam.sdk.io.synthetic.SyntheticSourceOptions.ProgressShape;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.SourceTestUtils;
import org.apache.beam.sdk.values.KV;
import org.apache.commons.math3.distribution.ConstantRealDistribution;
@@ -39,10 +40,9 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.ExpectedException;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Unit tests for {@link SyntheticBoundedSource}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class SyntheticBoundedSourceTest {
@Rule public final ExpectedException thrown = ExpectedException.none();
diff --git
a/sdks/java/io/synthetic/src/test/java/org/apache/beam/sdk/io/synthetic/SyntheticStepTest.java
b/sdks/java/io/synthetic/src/test/java/org/apache/beam/sdk/io/synthetic/SyntheticStepTest.java
index b400dd020a7..b4d65bd9e64 100644
---
a/sdks/java/io/synthetic/src/test/java/org/apache/beam/sdk/io/synthetic/SyntheticStepTest.java
+++
b/sdks/java/io/synthetic/src/test/java/org/apache/beam/sdk/io/synthetic/SyntheticStepTest.java
@@ -21,6 +21,7 @@ import static org.junit.Assert.assertEquals;
import java.nio.ByteBuffer;
import java.util.List;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
import org.apache.beam.sdk.transforms.Create;
@@ -34,10 +35,9 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.ExpectedException;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Unit tests for {@link SyntheticStep}. */
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class SyntheticStepTest {
@Rule public final ExpectedException thrown = ExpectedException.none();
@Rule public final transient TestPipeline p = TestPipeline.create();
diff --git a/sdks/java/testing/nexmark/build.gradle
b/sdks/java/testing/nexmark/build.gradle
index 0eeaf931a88..9dab21aae88 100644
--- a/sdks/java/testing/nexmark/build.gradle
+++ b/sdks/java/testing/nexmark/build.gradle
@@ -200,3 +200,8 @@ task run(type: JavaExec) {
classpath = configurations.gradleRun
args nexmarkArgsList.toArray()
}
+
+test {
+ maxParallelForks 4
+}
+
diff --git
a/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/BoundedSideInputJoinTest.java
b/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/BoundedSideInputJoinTest.java
index 9079c5497c2..5087fe999bf 100644
---
a/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/BoundedSideInputJoinTest.java
+++
b/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/BoundedSideInputJoinTest.java
@@ -30,6 +30,7 @@ import org.apache.beam.sdk.nexmark.NexmarkUtils;
import org.apache.beam.sdk.nexmark.model.Bid;
import org.apache.beam.sdk.nexmark.model.Event;
import org.apache.beam.sdk.nexmark.model.KnownSize;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -45,10 +46,9 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Test the various NEXMark queries yield results coherent with their models.
*/
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
@SuppressWarnings({
"rawtypes" // TODO(https://github.com/apache/beam/issues/20447)
})
diff --git
a/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/QueryTest.java
b/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/QueryTest.java
index 2039ea6c5b7..282f6b244da 100644
---
a/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/QueryTest.java
+++
b/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/QueryTest.java
@@ -22,6 +22,7 @@ import org.apache.beam.sdk.nexmark.NexmarkConfiguration;
import org.apache.beam.sdk.nexmark.NexmarkUtils;
import org.apache.beam.sdk.nexmark.model.Event;
import org.apache.beam.sdk.nexmark.model.KnownSize;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -34,10 +35,9 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Test the various NEXMark queries yield results coherent with their models.
*/
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
public class QueryTest {
private static final NexmarkConfiguration CONFIG =
NexmarkConfiguration.DEFAULT.copy();
diff --git
a/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/SessionSideInputJoinTest.java
b/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/SessionSideInputJoinTest.java
index 1659cccd73e..70625a58201 100644
---
a/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/SessionSideInputJoinTest.java
+++
b/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/SessionSideInputJoinTest.java
@@ -30,6 +30,7 @@ import org.apache.beam.sdk.nexmark.NexmarkUtils;
import org.apache.beam.sdk.nexmark.model.Bid;
import org.apache.beam.sdk.nexmark.model.Event;
import org.apache.beam.sdk.nexmark.model.KnownSize;
+import org.apache.beam.sdk.testing.BeamParallelJunit4Runner;
import org.apache.beam.sdk.testing.NeedsRunner;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
@@ -47,10 +48,9 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Test the various NEXMark queries yield results coherent with their models.
*/
-@RunWith(JUnit4.class)
+@RunWith(BeamParallelJunit4Runner.class)
@SuppressWarnings({
"rawtypes" // TODO(https://github.com/apache/beam/issues/20447)
})
diff --git
a/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/SqlQueryTest.java
b/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/SqlQueryTest.java
index 854149893dd..b7ea9a38339 100644
---
a/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/SqlQueryTest.java
+++
b/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/SqlQueryTest.java
@@ -39,7 +39,6 @@ import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.experimental.runners.Enclosed;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Test the various NEXMark queries yield results coherent with their models.
*/
@RunWith(Enclosed.class)
@@ -144,7 +143,7 @@ public class SqlQueryTest {
}
}
- @RunWith(JUnit4.class)
+ @RunWith(org.apache.beam.sdk.testing.BeamParallelJunit4Runner.class)
public static class SqlQueryTestCalcite extends SqlQueryTestCases {
@Override
protected SqlQuery1 getQuery1() {
diff --git
a/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/sql/SqlBoundedSideInputJoinTest.java
b/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/sql/SqlBoundedSideInputJoinTest.java
index 0a296f16442..952ce5f0831 100644
---
a/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/sql/SqlBoundedSideInputJoinTest.java
+++
b/sdks/java/testing/nexmark/src/test/java/org/apache/beam/sdk/nexmark/queries/sql/SqlBoundedSideInputJoinTest.java
@@ -49,7 +49,6 @@ import org.junit.Rule;
import org.junit.Test;
import org.junit.experimental.runners.Enclosed;
import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
/** Test the various NEXMark queries yield results coherent with their models.
*/
@RunWith(Enclosed.class)
@@ -213,7 +212,7 @@ public class SqlBoundedSideInputJoinTest {
}
}
- @RunWith(JUnit4.class)
+ @RunWith(org.apache.beam.sdk.testing.BeamParallelJunit4Runner.class)
public static class SqlBoundedSideInputJoinTestCalcite extends
SqlBoundedSideInputJoinTestCases {
@Override
protected SqlBoundedSideInputJoin getQuery(NexmarkConfiguration
configuration) {