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 fa34c16167c Make ValidatesRunner faster: Share and cache GCS staging 
across TestDataflowRunner executions (#40321)
fa34c16167c is described below

commit fa34c16167c9b587ad648d632d60a8bef27fe37c
Author: Reuven Lax <[email protected]>
AuthorDate: Tue Sep 29 12:46:22 2026 -0700

    Make ValidatesRunner faster: Share and cache GCS staging across 
TestDataflowRunner executions (#40321)
    
    * Share and cache GCS staging across TestDataflowRunner executions
    
    * fixes
    
    * address comments
---
 runners/google-cloud-dataflow-java/build.gradle    |  2 +
 .../beam/runners/dataflow/TestDataflowRunner.java  |  6 +++
 .../beam/runners/dataflow/util/GcsStager.java      | 45 ++++++++++++++++++++++
 .../runners/dataflow/TestDataflowRunnerTest.java   | 20 ++++++++++
 4 files changed, 73 insertions(+)

diff --git a/runners/google-cloud-dataflow-java/build.gradle 
b/runners/google-cloud-dataflow-java/build.gradle
index d52d2f1bc7b..af42dc2a501 100644
--- a/runners/google-cloud-dataflow-java/build.gradle
+++ b/runners/google-cloud-dataflow-java/build.gradle
@@ -173,6 +173,7 @@ def legacyPipelineOptions = [
   "--project=${gcpProject}",
   "--region=${gcpRegion}",
   "--tempRoot=${dataflowValidatesTempRoot}",
+  "--stagingLocation=${dataflowValidatesTempRoot}/staging",
   "--dataflowWorkerJar=${dataflowLegacyWorkerJar}",
   "--numWorkers=1",
   "--maxNumWorkers=1",
@@ -193,6 +194,7 @@ def runnerV2CommonPipelineOptions = [
   "--project=${gcpProject}",
   "--region=${gcpRegion}",
   "--tempRoot=${dataflowValidatesTempRoot}",
+  "--stagingLocation=${dataflowValidatesTempRoot}/staging",
   "--experiments=use_unified_worker,use_runner_v2",
   "--firestoreDb=${firestoreDb}",
   "--numWorkers=1",
diff --git 
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/TestDataflowRunner.java
 
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/TestDataflowRunner.java
index db8364bcbe8..89d198154ad 100644
--- 
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/TestDataflowRunner.java
+++ 
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/TestDataflowRunner.java
@@ -85,6 +85,12 @@ public class TestDataflowRunner extends 
PipelineRunner<DataflowPipelineJob> {
       tempLocation = tempLocation.substring(0, tempLocation.length() - 
File.separator.length());
     }
     dataflowOptions.setTempLocation(tempLocation);
+    String defaultPerJobStagingLocation =
+        FileSystems.matchNewDirectory(tempLocation, "staging").toString();
+    if 
(defaultPerJobStagingLocation.equals(dataflowOptions.getStagingLocation())) {
+      dataflowOptions.setStagingLocation(
+          FileSystems.matchNewDirectory(dataflowOptions.getTempRoot(), 
"staging").toString());
+    }
 
     return new TestDataflowRunner(
         dataflowOptions, 
DataflowClient.create(options.as(DataflowPipelineOptions.class)));
diff --git 
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/GcsStager.java
 
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/GcsStager.java
index bf34e007c40..413a870a365 100644
--- 
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/GcsStager.java
+++ 
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/GcsStager.java
@@ -21,15 +21,42 @@ import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Mo
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
 
 import com.google.api.services.dataflow.model.DataflowPackage;
+import com.google.auto.value.AutoValue;
+import java.time.Duration;
+import java.util.Collections;
 import java.util.List;
+import java.util.concurrent.ExecutionException;
 import org.apache.beam.runners.dataflow.options.DataflowPipelineOptions;
 import org.apache.beam.runners.dataflow.util.PackageUtil.StagedFile;
 import org.apache.beam.sdk.extensions.gcp.storage.GcsCreateOptions;
 import org.apache.beam.sdk.options.PipelineOptions;
 import org.apache.beam.sdk.util.MimeTypes;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Throwables;
+import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.Cache;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.CacheBuilder;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.UncheckedExecutionException;
 
 /** Utility class for staging files to GCS. */
 public class GcsStager implements Stager {
+  @AutoValue
+  abstract static class StagedFilesCacheKey {
+    abstract String getStagingLocation();
+
+    abstract List<StagedFile> getFilesToStage();
+
+    static StagedFilesCacheKey of(String stagingLocation, List<StagedFile> 
filesToStage) {
+      return new AutoValue_GcsStager_StagedFilesCacheKey(stagingLocation, 
filesToStage);
+    }
+  }
+
+  private static final int MAX_STAGED_FILES_CACHE_SIZE = 5000;
+
+  private static final Cache<StagedFilesCacheKey, List<DataflowPackage>> 
STAGED_FILES_CACHE =
+      CacheBuilder.newBuilder()
+          .maximumSize(MAX_STAGED_FILES_CACHE_SIZE)
+          .expireAfterWrite(Duration.ofMinutes(30))
+          .build();
+
   private DataflowPipelineOptions options;
 
   private GcsStager(DataflowPipelineOptions options) {
@@ -49,6 +76,24 @@ public class GcsStager implements Stager {
    */
   @Override
   public List<DataflowPackage> stageFiles(List<StagedFile> filesToStage) {
+    String stagingLocation = options.getStagingLocation();
+    if (stagingLocation != null) {
+      StagedFilesCacheKey cacheKey = StagedFilesCacheKey.of(stagingLocation, 
filesToStage);
+      try {
+        return STAGED_FILES_CACHE.get(
+            cacheKey, () -> 
Collections.unmodifiableList(stageFilesUncached(filesToStage)));
+      } catch (ExecutionException | UncheckedExecutionException e) {
+        if (e.getCause() != null) {
+          Throwables.throwIfUnchecked(e.getCause());
+          throw new RuntimeException(e.getCause());
+        }
+        throw new RuntimeException(e);
+      }
+    }
+    return stageFilesUncached(filesToStage);
+  }
+
+  private List<DataflowPackage> stageFilesUncached(List<StagedFile> 
filesToStage) {
     try (PackageUtil packageUtil = PackageUtil.withDefaultThreadPool()) {
       return packageUtil.stageClasspathElements(
           filesToStage, options.getStagingLocation(), buildCreateOptions());
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 ed6259a3ee2..4f6ff01c327 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
@@ -49,6 +49,7 @@ import org.apache.beam.sdk.PipelineResult.State;
 import org.apache.beam.sdk.extensions.gcp.auth.TestCredential;
 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.PAssert;
 import org.apache.beam.sdk.testing.SerializableMatcher;
@@ -94,6 +95,7 @@ public class TestDataflowRunnerTest {
     options.setGcpCredential(new TestCredential());
     options.setRunner(TestDataflowRunner.class);
     options.setPathValidatorClass(NoopPathValidator.class);
+    FileSystems.setDefaultPipelineOptions(options);
   }
 
   @Test
@@ -102,6 +104,24 @@ public class TestDataflowRunnerTest {
         "TestDataflowRunner#TestAppName", 
TestDataflowRunner.fromOptions(options).toString());
   }
 
+  @Test
+  public void testFromOptionsUsesSharedStagingLocationUnderTempRoot() {
+    options.setJobName("test-job-1");
+    TestDataflowRunner.fromOptions(options);
+    assertEquals("gs://test/test-job-1/output/results", 
options.getTempLocation());
+    assertEquals("gs://test/test-job-1/output/results", 
options.getGcpTempLocation());
+    assertEquals("gs://test/staging/", options.getStagingLocation());
+  }
+
+  @Test
+  public void testFromOptionsPreservesExplicitStagingLocation() {
+    options.setJobName("test-job-2");
+    options.setStagingLocation("gs://custom-bucket/custom-staging/");
+    TestDataflowRunner.fromOptions(options);
+    assertEquals("gs://test/test-job-2/output/results", 
options.getTempLocation());
+    assertEquals("gs://custom-bucket/custom-staging/", 
options.getStagingLocation());
+  }
+
   @Test
   public void testRunBatchJobThatSucceeds() throws Exception {
     Pipeline p = Pipeline.create(options);

Reply via email to