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