Abacn commented on code in PR #40320:
URL: https://github.com/apache/beam/pull/40320#discussion_r4126717738
##########
sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java:
##########
@@ -98,6 +104,11 @@ public class Environments {
.put(ENVIRONMENT_PROCESS, ImmutableSet.of(processCommandOption,
processVariablesOption))
.build();
+ private static final ConcurrentHashMap<FileHashCacheKey, HashCode>
FILE_HASH_CACHE =
Review Comment:
In general, consider a bounded size cache in case of long running JVMs. One
can use Guava's com.google.common.cache.Cache with maximumSize. Same applies to
#40321
##########
sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java:
##########
@@ -538,6 +549,58 @@ public static String createStagingFileName(File path,
HashCode hash) {
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.computeIfAbsent(
+ key,
+ k -> {
+ try {
+ return Files.asByteSource(file).hash(Hashing.sha256());
+ } catch (IOException e) {
+ throw new UncheckedIOException(e);
+ }
+ });
+ } catch (UncheckedIOException e) {
+ throw 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
Review Comment:
Technically this isn't a true "equal" but just for file of same path,
length, and same modified time, it will be treated equal.
We had a couple of security reports about file staging recently not sure if
this would have security implications
Again, setting a retention time for the cache may mitigate potential
conflicts
##########
sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java:
##########
@@ -549,11 +612,52 @@ public static String
getExternalServiceAddress(PortablePipelineOptions options)
}
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");
Review Comment:
(Note: this fixed a pre-existing temp file leak which is good)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]