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

Reply via email to