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);
     }

Reply via email to