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 bb9d98aec5f Merge pull request #40323 from
reuvenlax/improve_test_tune_validates_runner
bb9d98aec5f is described below
commit bb9d98aec5f6caa4c22d6db988b2caf132240ae7
Author: Reuven Lax <[email protected]>
AuthorDate: Tue Sep 29 08:54:01 2026 -0700
Merge pull request #40323 from reuvenlax/improve_test_tune_validates_runner
Make ValidatesRunner faster: Tune Dataflow ValidatesRunner configuration
---
runners/google-cloud-dataflow-java/build.gradle | 29 ++++++++++++++++++++--
.../apache/beam/sdk/transforms/GroupByKeyTest.java | 2 --
.../org/apache/beam/sdk/transforms/ParDoTest.java | 17 ++++++++++---
3 files changed, 41 insertions(+), 7 deletions(-)
diff --git a/runners/google-cloud-dataflow-java/build.gradle
b/runners/google-cloud-dataflow-java/build.gradle
index fbfa837d9fd..d52d2f1bc7b 100644
--- a/runners/google-cloud-dataflow-java/build.gradle
+++ b/runners/google-cloud-dataflow-java/build.gradle
@@ -174,6 +174,9 @@ def legacyPipelineOptions = [
"--region=${gcpRegion}",
"--tempRoot=${dataflowValidatesTempRoot}",
"--dataflowWorkerJar=${dataflowLegacyWorkerJar}",
+ "--numWorkers=1",
+ "--maxNumWorkers=1",
+ "--diskSizeGb=30",
"--usePublicIps=false",
"--experiments=enable_lineage"
]
@@ -192,6 +195,9 @@ def runnerV2CommonPipelineOptions = [
"--tempRoot=${dataflowValidatesTempRoot}",
"--experiments=use_unified_worker,use_runner_v2",
"--firestoreDb=${firestoreDb}",
+ "--numWorkers=1",
+ "--maxNumWorkers=1",
+ "--diskSizeGb=30",
"--usePublicIps=false",
"--experiments=enable_lineage"
]
@@ -228,6 +234,18 @@ def commonRunnerV2ExcludeCategories = [
'org.apache.beam.sdk.testing.UsesBoundedTrieMetrics', // Dataflow QM as of
now does not support returning back BoundedTrie in metric result.
]
+def isValidatesRunnerTestClass = { FileTreeElement element ->
+ if (element.isDirectory()) {
+ return true
+ }
+ if (!element.name.endsWith('.class') ||
+ element.name.startsWith('ValidateRunnerXlangTest')) {
+ return false
+ }
+ return new String(element.file.bytes,
java.nio.charset.StandardCharsets.ISO_8859_1)
+ .contains('Lorg/apache/beam/sdk/testing/ValidatesRunner;')
+}
+
def createLegacyWorkerValidatesRunnerTest = { Map args ->
def name = args.name
def pipelineOptions = args.pipelineOptions ?: legacyPipelineOptions
@@ -245,6 +263,7 @@ def createLegacyWorkerValidatesRunnerTest = { Map args ->
classpath = configurations.validatesRunner
testClassesDirs =
files(project(":sdks:java:core").sourceSets.test.output.classesDirs) +
files(project(project.path).sourceSets.test.output.classesDirs)
+ include isValidatesRunnerTestClass
useJUnit {
includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner'
commonLegacyExcludeCategories.each {
@@ -278,6 +297,7 @@ def createRunnerV2ValidatesRunnerTest = { Map args ->
classpath = configurations.validatesRunner
testClassesDirs =
files(project(":sdks:java:core").sourceSets.test.output.classesDirs) +
files(project(project.path).sourceSets.test.output.classesDirs)
+ include isValidatesRunnerTestClass
useJUnit {
includeCategories 'org.apache.beam.sdk.testing.ValidatesRunner'
commonRunnerV2ExcludeCategories.each {
@@ -466,8 +486,13 @@ task validatesRunner {
'org.apache.beam.sdk.transforms.ParDoLifecycleTest.testTeardownCalledAfterExceptionInStartBundle',
'org.apache.beam.sdk.transforms.ParDoLifecycleTest.testTeardownCalledAfterExceptionInStartBundleStateful',
],
- // Batch legacy worker does not support bundle finalization.
- excludedCategories: [ 'org.apache.beam.sdk.testing.UsesBundleFinalizer', ],
+ // Batch legacy worker does not support bundle finalization, triggered
side inputs, or unbounded PCollections.
+ excludedCategories: [
+ 'org.apache.beam.sdk.testing.UsesBundleFinalizer',
+ 'org.apache.beam.sdk.testing.UsesTriggeredSideInputs',
+ 'org.apache.beam.sdk.testing.UsesUnboundedPCollections',
+ 'org.apache.beam.sdk.testing.UsesUnboundedSplittableParDo',
+ ],
))
}
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/GroupByKeyTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/GroupByKeyTest.java
index 5464838ad4d..18541437d5f 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/GroupByKeyTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/GroupByKeyTest.java
@@ -94,7 +94,6 @@ import org.junit.Assert;
import org.junit.Rule;
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;
@@ -104,7 +103,6 @@ import org.junit.runners.JUnit4;
"unchecked",
"unused"
})
-@RunWith(Enclosed.class)
public class GroupByKeyTest implements Serializable {
/** Shared test base class with setup/teardown helpers. */
public abstract static class SharedTestBase {
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ParDoTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ParDoTest.java
index 6beea338689..7eb704b2bbf 100644
--- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ParDoTest.java
+++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/ParDoTest.java
@@ -3356,7 +3356,8 @@ public class ParDoTest implements Serializable {
ValidatesRunner.class,
UsesStatefulParDo.class,
UsesOrderedListState.class,
- UsesOnWindowExpiration.class
+ UsesOnWindowExpiration.class,
+ UsesUnboundedPCollections.class
})
public void testOrderedListStateUnbounded() {
testOrderedListStateImpl(true);
@@ -3420,7 +3421,12 @@ public class ParDoTest implements Serializable {
}
@Test
- @Category({ValidatesRunner.class, UsesStatefulParDo.class,
UsesOrderedListState.class})
+ @Category({
+ ValidatesRunner.class,
+ UsesStatefulParDo.class,
+ UsesOrderedListState.class,
+ UsesUnboundedPCollections.class
+ })
public void testOrderedListStateRangeFetchUnbounded() {
testOrderedListStateRangeFetchImpl(true);
}
@@ -3493,7 +3499,12 @@ public class ParDoTest implements Serializable {
}
@Test
- @Category({ValidatesRunner.class, UsesStatefulParDo.class,
UsesOrderedListState.class})
+ @Category({
+ ValidatesRunner.class,
+ UsesStatefulParDo.class,
+ UsesOrderedListState.class,
+ UsesUnboundedPCollections.class
+ })
public void testOrderedListStateRangeDeleteUnbounded() {
testOrderedListStateRangeDeleteImpl(true);
}