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 01b18a16867 Make ValidatesRunner faster: Improvements to file staging
(#40320)
01b18a16867 is described below
commit 01b18a16867046a5ded5c02c5b703e1879d0dd3d
Author: Reuven Lax <[email protected]>
AuthorDate: Tue Sep 29 10:34:42 2026 -0700
Make ValidatesRunner faster: Improvements to file staging (#40320)
* Cache classpath resource detection, artifact file hashes, and directory
zips
* update classpath code
* fixes
---
.../beam/runners/dataflow/util/PackageUtil.java | 4 +-
.../beam/sdk/util/construction/Environments.java | 114 +++++++++++++++++++--
.../ClasspathScanningResourcesDetector.java | 64 ++++++++++--
.../sdk/util/construction/EnvironmentsTest.java | 53 ++++++++++
4 files changed, 217 insertions(+), 18 deletions(-)
diff --git
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/PackageUtil.java
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/PackageUtil.java
index 8124c74deec..360fd06d7bd 100644
---
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/PackageUtil.java
+++
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/PackageUtil.java
@@ -425,9 +425,7 @@ public class PackageUtil implements Closeable {
switch (dest) {
case "dataflow-worker.jar":
case "windmill_main":
- target =
- Environments.createStagingFileName(
- file, Files.asByteSource(file).hash(Hashing.sha256()));
+ target = Environments.createStagingFileName(file,
Environments.getFileHash(file));
LOG.info("Staging custom {} as {}", dest, target);
break;
default:
diff --git
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java
index d46d85e63ef..c02de598bce 100644
---
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java
+++
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java
@@ -22,12 +22,18 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import java.io.File;
import java.io.FileOutputStream;
import java.io.IOException;
+import java.io.UncheckedIOException;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.UUID;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import org.apache.beam.model.pipeline.v1.Endpoints.ApiServiceDescriptor;
@@ -55,10 +61,13 @@ import
org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.InvalidProtocolBu
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Strings;
+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.collect.ImmutableList;
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.ImmutableSet;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.HashCode;
+import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.Hasher;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.Hashing;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.io.Files;
import org.checkerframework.checker.nullness.qual.Nullable;
@@ -98,6 +107,11 @@ public class Environments {
.put(ENVIRONMENT_PROCESS, ImmutableSet.of(processCommandOption,
processVariablesOption))
.build();
+ private static final Cache<FileHashCacheKey, HashCode> FILE_HASH_CACHE =
+ CacheBuilder.newBuilder().maximumSize(10_000).build();
+ private static final ConcurrentHashMap<HashCode, File> DIRECTORY_ZIP_CACHE =
+ new ConcurrentHashMap<>();
+
public enum JavaVersion {
java11("java11", "11", 11),
java17("java17", "17", 17),
@@ -434,7 +448,7 @@ public class Environments {
File zippedFile;
try {
zippedFile = zipDirectory(file);
- hashCode = Files.asByteSource(zippedFile).hash(Hashing.sha256());
+ hashCode = getFileHash(zippedFile);
} catch (IOException e) {
throw new RuntimeException(e);
}
@@ -448,7 +462,7 @@ public class Environments {
} else {
try {
- hashCode = Files.asByteSource(file).hash(Hashing.sha256());
+ hashCode = getFileHash(file);
} catch (IOException e) {
throw new RuntimeException(e);
}
@@ -538,6 +552,53 @@ public class Environments {
return String.format("%s-%s%s", fileName, encodedHash, suffix);
}
+ /**
+ * Returns the SHA-256 {@link HashCode} for {@code file}, caching the result
in memory keyed by
+ * the file's absolute path, length, and last-modified timestamp.
+ */
+ public static HashCode getFileHash(File file) throws IOException {
+ FileHashCacheKey key = new FileHashCacheKey(file);
+ try {
+ return FILE_HASH_CACHE.get(key, () ->
Files.asByteSource(file).hash(Hashing.sha256()));
+ } catch (ExecutionException e) {
+ if (e.getCause() instanceof IOException) {
+ throw (IOException) e.getCause();
+ }
+ throw new RuntimeException(e.getCause());
+ }
+ }
+
+ private static final class FileHashCacheKey {
+ private final String absolutePath;
+ private final long length;
+ private final long lastModified;
+
+ FileHashCacheKey(File file) {
+ this.absolutePath = file.getAbsolutePath();
+ this.length = file.length();
+ this.lastModified = file.lastModified();
+ }
+
+ @Override
+ public boolean equals(@Nullable Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (!(o instanceof FileHashCacheKey)) {
+ return false;
+ }
+ FileHashCacheKey that = (FileHashCacheKey) o;
+ return length == that.length
+ && lastModified == that.lastModified
+ && Objects.equals(absolutePath, that.absolutePath);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(absolutePath, length, lastModified);
+ }
+ }
+
public static String getExternalServiceAddress(PortablePipelineOptions
options) {
String environmentConfig = options.getDefaultEnvironmentConfig();
String environmentOption =
@@ -549,11 +610,52 @@ public class Environments {
}
private static File zipDirectory(File directory) throws IOException {
- File zipFile = File.createTempFile(directory.getName(), ".zip");
- try (FileOutputStream fos = new FileOutputStream(zipFile)) {
- ZipFiles.zipDirectory(directory, fos);
+ HashCode metadataHash = computeDirectoryMetadataHash(directory);
+ try {
+ return DIRECTORY_ZIP_CACHE.compute(
+ metadataHash,
+ (k, cachedZip) -> {
+ if (cachedZip != null && cachedZip.exists()) {
+ return cachedZip;
+ }
+ try {
+ File zipFile = File.createTempFile(directory.getName(), ".zip");
+ zipFile.deleteOnExit();
+ try (FileOutputStream fos = new FileOutputStream(zipFile)) {
+ ZipFiles.zipDirectory(directory, fos);
+ }
+ return zipFile;
+ } catch (IOException e) {
+ throw new UncheckedIOException(e);
+ }
+ });
+ } catch (UncheckedIOException e) {
+ throw e.getCause();
+ }
+ }
+
+ private static HashCode computeDirectoryMetadataHash(File directory) {
+ Hasher hasher = Hashing.sha256().newHasher();
+ hasher.putString(directory.getAbsolutePath(), StandardCharsets.UTF_8);
+ hashDirectoryMetadataRecursive(directory, "", hasher);
+ return hasher.hash();
+ }
+
+ private static void hashDirectoryMetadataRecursive(File inputFile, String
prefix, Hasher hasher) {
+ String entryName = prefix + inputFile.getName();
+ hasher.putString(entryName, StandardCharsets.UTF_8);
+ if (inputFile.isDirectory()) {
+ File[] children = inputFile.listFiles();
+ if (children != null) {
+ Arrays.sort(children);
+ for (File child : children) {
+ hashDirectoryMetadataRecursive(child, entryName + "/", hasher);
+ }
+ }
+ } else {
+ hasher.putLong(inputFile.length());
+ hasher.putLong(inputFile.lastModified());
}
- return zipFile;
}
private static class ProcessPayloadReferenceJSON {
diff --git
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/resources/ClasspathScanningResourcesDetector.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/resources/ClasspathScanningResourcesDetector.java
index 4c289603ae0..5a5c6b32527 100644
---
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/resources/ClasspathScanningResourcesDetector.java
+++
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/resources/ClasspathScanningResourcesDetector.java
@@ -19,8 +19,13 @@ package org.apache.beam.sdk.util.construction.resources;
import io.github.classgraph.ClassGraph;
import java.io.File;
+import java.lang.ref.WeakReference;
+import java.util.ArrayList;
+import java.util.Collections;
import java.util.List;
+import java.util.Objects;
import java.util.stream.Collectors;
+import org.checkerframework.checker.nullness.qual.Nullable;
/**
* Attempts to detect all the resources to be staged using classgraph library.
@@ -30,6 +35,27 @@ import java.util.stream.Collectors;
*/
public class ClasspathScanningResourcesDetector implements
PipelineResourcesDetector {
+ private static final class CachedClasspath {
+ private final WeakReference<ClassLoader> classLoader;
+ private final @Nullable String javaClassPath;
+ private final List<String> files;
+
+ CachedClasspath(ClassLoader classLoader, @Nullable String javaClassPath,
List<String> files) {
+ this.classLoader = new WeakReference<>(classLoader);
+ this.javaClassPath = javaClassPath;
+ this.files = Collections.unmodifiableList(new ArrayList<>(files));
+ }
+
+ boolean matches(@Nullable ClassLoader loader, @Nullable String
currentJavaClassPath) {
+ return loader != null
+ && classLoader.get() == loader
+ && Objects.equals(javaClassPath, currentJavaClassPath);
+ }
+ }
+
+ private static final Object LOCK = new Object();
+ private static volatile @Nullable CachedClasspath cachedClasspath;
+
private transient ClassGraph classGraph;
public ClasspathScanningResourcesDetector(ClassGraph classGraph) {
@@ -43,14 +69,34 @@ public class ClasspathScanningResourcesDetector implements
PipelineResourcesDete
* @return A list of absolute paths to the resources the class loader uses.
*/
@Override
- public List<String> detect(ClassLoader classLoader) {
- List<File> classpathContents =
- classGraph
- .disableNestedJarScanning()
- .addClassLoader(classLoader)
- .scan(1)
- .getClasspathFiles();
-
- return
classpathContents.stream().map(File::getAbsolutePath).collect(Collectors.toList());
+ public List<String> detect(@Nullable ClassLoader classLoader) {
+ String currentJavaClassPath = System.getProperty("java.class.path");
+ CachedClasspath snapshot = cachedClasspath;
+ if (snapshot != null && snapshot.matches(classLoader,
currentJavaClassPath)) {
+ return new ArrayList<>(snapshot.files);
+ }
+
+ synchronized (LOCK) {
+ currentJavaClassPath = System.getProperty("java.class.path");
+ snapshot = cachedClasspath;
+ if (snapshot != null && snapshot.matches(classLoader,
currentJavaClassPath)) {
+ return new ArrayList<>(snapshot.files);
+ }
+
+ List<File> classpathContents;
+ if (classLoader != null) {
+ classpathContents =
+
classGraph.disableNestedJarScanning().addClassLoader(classLoader).getClasspathFiles();
+ } else {
+ classpathContents =
classGraph.disableNestedJarScanning().getClasspathFiles();
+ }
+
+ List<String> result =
+
classpathContents.stream().map(File::getAbsolutePath).collect(Collectors.toList());
+ if (classLoader != null) {
+ cachedClasspath = new CachedClasspath(classLoader,
currentJavaClassPath, result);
+ }
+ return result;
+ }
}
}
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/construction/EnvironmentsTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/construction/EnvironmentsTest.java
index ebd4e9fbe24..b4b71a43d00 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/construction/EnvironmentsTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/construction/EnvironmentsTest.java
@@ -24,12 +24,15 @@ import static org.hamcrest.Matchers.equalTo;
import static org.hamcrest.Matchers.hasItem;
import static org.hamcrest.Matchers.hasSize;
import static org.hamcrest.Matchers.is;
+import static org.hamcrest.Matchers.not;
import static org.hamcrest.Matchers.notNullValue;
import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertSame;
import java.io.File;
import java.io.IOException;
import java.io.Serializable;
+import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.List;
import java.util.Optional;
@@ -58,6 +61,8 @@ import org.apache.beam.sdk.values.TupleTag;
import org.apache.beam.sdk.values.TupleTagList;
import org.apache.beam.sdk.values.WindowingStrategy;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.HashCode;
+import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.io.Files;
import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.ExpectedException;
@@ -399,4 +404,52 @@ public class EnvironmentsTest implements Serializable {
env, BeamUrns.getUrn(StandardEnvironments.Environments.EXTERNAL)),
notNullValue());
}
+
+ @Test
+ public void testGetFileHashCachingAndInvalidation() throws Exception {
+ File file = File.createTempFile("hash-cache-test-", ".txt");
+ file.deleteOnExit();
+ Files.asCharSink(file, StandardCharsets.UTF_8).write("initial-content");
+
+ HashCode hash1 = Environments.getFileHash(file);
+ HashCode hash2 = Environments.getFileHash(file);
+ assertSame(hash1, hash2);
+
+ // Modifying the file content and length invalidates the cache entry.
+ Files.asCharSink(file,
StandardCharsets.UTF_8).write("updated-longer-content");
+ HashCode hash3 = Environments.getFileHash(file);
+ assertThat(hash3, not(equalTo(hash1)));
+ }
+
+ @Test
+ public void testGetArtifactsDirectoryZipCachingAndInvalidation() throws
Exception {
+ File tempDir = Files.createTempDir();
+ tempDir.deleteOnExit();
+ File child = new File(tempDir, "entry.txt");
+ child.deleteOnExit();
+ Files.asCharSink(child, StandardCharsets.UTF_8).write("v1");
+
+ List<ArtifactInformation> firstArtifacts =
+ Environments.getArtifacts(ImmutableList.of(tempDir.getAbsolutePath()));
+ List<ArtifactInformation> secondArtifacts =
+ Environments.getArtifacts(ImmutableList.of(tempDir.getAbsolutePath()));
+
+ RunnerApi.ArtifactFilePayload firstPayload =
+
RunnerApi.ArtifactFilePayload.parseFrom(firstArtifacts.get(0).getTypePayload());
+ RunnerApi.ArtifactFilePayload secondPayload =
+
RunnerApi.ArtifactFilePayload.parseFrom(secondArtifacts.get(0).getTypePayload());
+
+ // Unchanged directory reuses the same cached zip file and SHA-256 hash.
+ assertEquals(firstPayload.getPath(), secondPayload.getPath());
+ assertEquals(firstPayload.getSha256(), secondPayload.getSha256());
+
+ // Modifying a file inside the directory produces a new zip and hash.
+ Files.asCharSink(child, StandardCharsets.UTF_8).write("v2-modified");
+ List<ArtifactInformation> thirdArtifacts =
+ Environments.getArtifacts(ImmutableList.of(tempDir.getAbsolutePath()));
+ RunnerApi.ArtifactFilePayload thirdPayload =
+
RunnerApi.ArtifactFilePayload.parseFrom(thirdArtifacts.get(0).getTypePayload());
+ assertThat(thirdPayload.getPath(), not(equalTo(firstPayload.getPath())));
+ assertThat(thirdPayload.getSha256(),
not(equalTo(firstPayload.getSha256())));
+ }
}