This is an automated email from the ASF dual-hosted git repository.
shunping 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 e26d6bacb93 Parallelize per-file rewrites in GcsUtilV2 copy and move
(#40376)
e26d6bacb93 is described below
commit e26d6bacb937561d40975788415d72fde17c8678
Author: Shunping Huang <[email protected]>
AuthorDate: Wed Oct 7 13:53:39 2026 -0400
Parallelize per-file rewrites in GcsUtilV2 copy and move (#40376)
* Parallelize per-file rewrites in GcsUtilV2 copy and move
* Add gcsMaxConcurrentRewrites option for GcsUtilV2
- Make the number of concurrent copy/move rewrites in GcsUtilV2 configurable
via GcsOptions (defaults to 32).
- Deprecate gcsRewriteDataOpBatchLimit,
which only applies to GcsUtilV1.
* Wait for in-flight rewrites in GcsUtilV2 copy and move on failure
When one rewrite fails, files not yet started are cancelled, but the
executor's shutdownNow() also interrupted rewrites already in flight, and
the method could return while they were still running. Wait up to 5 minutes.
---
.../sdk/extensions/gcp/options/GcsOptions.java | 31 ++-
.../beam/sdk/extensions/gcp/util/GcsUtilV2.java | 244 ++++++++++++++++-----
.../sdk/extensions/gcp/util/GcsUtilV2Test.java | 75 +++++++
3 files changed, 294 insertions(+), 56 deletions(-)
diff --git
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptions.java
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptions.java
index fcb1587a120..16acb90d11a 100644
---
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptions.java
+++
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptions.java
@@ -171,11 +171,40 @@ public interface GcsOptions extends
ApplicationNameOptions, GcpOptions, Pipeline
void setGcsHttpRequestWriteTimeout(@Nullable Integer timeoutMs);
- @Description("Batching limit for rewrite ops which will copy data.")
+ /**
+ * Batching limit for rewrite ops which will copy data.
+ *
+ * @deprecated Only applicable to GcsUtil V1. Not applicable to GcsUtil V2
(enabled with the
+ * {@code use_gcsutil_v2} experiment), which does not batch rewrites;
use {@link
+ * #getGcsMaxConcurrentRewrites} to tune copy/move concurrency there
instead.
+ */
+ @Deprecated
+ @Description(
+ "Batching limit for rewrite ops which will copy data. Deprecated: NOT
applicable to GcsUtil"
+ + " V2 (enabled with the use_gcsutil_v2 experiment); use
gcsMaxConcurrentRewrites"
+ + " instead.")
@Nullable Integer getGcsRewriteDataOpBatchLimit();
+ /**
+ * @deprecated Only applicable to GcsUtil V1. See {@link
#getGcsRewriteDataOpBatchLimit}.
+ */
+ @Deprecated
void setGcsRewriteDataOpBatchLimit(@Nullable Integer timeoutMs);
+ /**
+ * Maximum number of files rewritten concurrently when copying or moving
files. If unset, a
+ * default of 32 is used.
+ *
+ * <p>Only applicable to GcsUtil V2, which is enabled with the {@code
use_gcsutil_v2} experiment.
+ */
+ @Description(
+ "Maximum number of files rewritten concurrently when copying or moving
files in GCS."
+ + " If unset, defaults to 32. Only applicable to GcsUtil V2 (enabled
with the"
+ + " use_gcsutil_v2 experiment).")
+ @Nullable Integer getGcsMaxConcurrentRewrites();
+
+ void setGcsMaxConcurrentRewrites(@Nullable Integer maxConcurrentRewrites);
+
/** If true, reports number of bytes written to each gcs bucket. */
@Description("Whether to report number of bytes written per GCS bucket.")
@Default.Boolean(false)
diff --git
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
index 20260aa3c79..dc52726fe17 100644
---
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
+++
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
@@ -52,6 +52,7 @@ import com.google.cloud.storage.StorageBatchResult;
import com.google.cloud.storage.StorageChannelUtils;
import com.google.cloud.storage.StorageException;
import com.google.cloud.storage.StorageOptions;
+import java.io.Closeable;
import java.io.FileNotFoundException;
import java.io.IOException;
import java.net.MalformedURLException;
@@ -67,6 +68,12 @@ import java.util.HashMap;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
+import java.util.concurrent.CancellationException;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
import java.util.regex.Pattern;
import org.apache.beam.runners.core.metrics.GcpResourceIdentifiers;
@@ -87,6 +94,7 @@ import org.apache.beam.sdk.options.PipelineOptions;
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.collect.ImmutableList;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.ThreadFactoryBuilder;
import org.checkerframework.checker.nullness.qual.Nullable;
class GcsUtilV2 {
@@ -141,6 +149,21 @@ class GcsUtilV2 {
*/
private static final long MEGABYTES_COPIED_PER_CHUNK = 2048L;
+ /**
+ * Default maximum number of files rewritten concurrently by {@link #copy}
and {@link #move}, used
+ * when {@link GcsOptions#getGcsMaxConcurrentRewrites} is unset.
+ */
+ private static final int DEFAULT_MAX_CONCURRENT_REWRITES = 32;
+
+ /**
+ * How long {@link #copy} and {@link #move} wait for in-flight rewrites to
finish after a failure
+ * before interrupting them.
+ */
+ private static final long REWRITE_TERMINATION_TIMEOUT_MINUTES = 5;
+
+ /** Maximum number of files rewritten concurrently by {@link #copy} and
{@link #move}. */
+ private final int maxConcurrentRewrites;
+
GcsUtilV2(PipelineOptions options) {
GcsOptions gcsOptions = options.as(GcsOptions.class);
this.projectId = options.as(GcpOptions.class).getProject();
@@ -173,6 +196,16 @@ class GcsUtilV2 {
? gcsOptions.getGcsWriteCounterPrefix()
: null);
this.gcsPerformanceMetrics =
Boolean.TRUE.equals(gcsOptions.getGcsPerformanceMetrics());
+
+ Integer configuredMaxConcurrentRewrites =
gcsOptions.getGcsMaxConcurrentRewrites();
+ checkArgument(
+ configuredMaxConcurrentRewrites == null ||
configuredMaxConcurrentRewrites > 0,
+ "gcsMaxConcurrentRewrites must be positive, but was %s",
+ configuredMaxConcurrentRewrites);
+ this.maxConcurrentRewrites =
+ configuredMaxConcurrentRewrites != null
+ ? configuredMaxConcurrentRewrites
+ : DEFAULT_MAX_CONCURRENT_REWRITES;
}
/**
@@ -594,73 +627,174 @@ class GcsUtilV2 {
srcList.size(),
dstList.size());
- for (int i = 0; i < srcList.size(); i++) {
- GcsPath srcPath = srcList.get(i);
- GcsPath dstPath = dstList.get(i);
- BlobId srcId = BlobId.of(srcPath.getBucket(), srcPath.getObject());
- BlobId dstId = BlobId.of(dstPath.getBucket(), dstPath.getObject());
+ if (srcList.isEmpty()) {
+ return;
+ }
+ if (srcList.size() == 1) {
+ // Nothing to overlap, so skip the thread pool.
+ rewriteOne(srcList.get(0), dstList.get(0), deleteSrc, srcMissing,
dstOverwrite);
+ return;
+ }
- CopyRequest.Builder copyRequestBuilder =
- CopyRequest.newBuilder()
- .setSource(srcId)
- .setMegabytesCopiedPerChunk(MEGABYTES_COPIED_PER_CHUNK);
+ // Each file costs up to three dependent round trips (target lookup,
rewrite, source delete),
+ // and java-storage cannot batch rewrites, so issue the files concurrently
instead.
+ int numThreads = Math.min(srcList.size(), maxConcurrentRewrites);
+ ExecutorService executor =
+ Executors.newFixedThreadPool(
+ numThreads,
+ new ThreadFactoryBuilder()
+ .setDaemon(true)
+ .setNameFormat("gcsutil-v2-rewrite-%d")
+ .build());
+ // Requests run on the pool's threads, so carry the caller's container
over to them.
+ MetricsContainer container = MetricsEnvironment.getCurrentContainer();
+ List<Future<Void>> futures = new ArrayList<>(srcList.size());
+ try {
+ for (int i = 0; i < srcList.size(); i++) {
+ GcsPath srcPath = srcList.get(i);
+ GcsPath dstPath = dstList.get(i);
+ futures.add(
+ executor.submit(
+ () -> {
+ if (container != null) {
+ try (Closeable scope =
MetricsEnvironment.scopedMetricsContainer(container)) {
+ rewriteOne(srcPath, dstPath, deleteSrc, srcMissing,
dstOverwrite);
+ }
+ } else {
+ rewriteOne(srcPath, dstPath, deleteSrc, srcMissing,
dstOverwrite);
+ }
+ return null;
+ }));
+ }
- if (dstOverwrite == OverwriteStrategy.ALWAYS_OVERWRITE) {
- copyRequestBuilder.setTarget(dstId);
- } else {
- // FAIL_IF_EXISTS, SKIP_IF_EXISTS and SAFE_OVERWRITE require checking
the target blob
- BlobInfo existingTarget;
+ // Like the sequential loop this replaces, stop at the first failure:
files not yet started
+ // are cancelled, while the ones already in flight are allowed to finish
(awaited below).
+ @Nullable Throwable failure = null;
+ for (Future<Void> future : futures) {
try {
- existingTarget = storage().get(dstId);
- } catch (StorageException e) {
- throw translateStorageException(dstPath, e);
- }
-
- if (existingTarget == null) {
- copyRequestBuilder.setTarget(dstId,
Storage.BlobTargetOption.doesNotExist());
- } else {
- switch (dstOverwrite) {
- case SKIP_IF_EXISTS:
- LOG.warn("Ignoring rewriting from {} to {} because target
exists.", srcPath, dstPath);
- continue; // Skip to next file in for-loop
-
- case SAFE_OVERWRITE:
- copyRequestBuilder.setTarget(
- dstId,
Storage.BlobTargetOption.generationMatch(existingTarget.getGeneration()));
- break;
-
- case FAIL_IF_EXISTS:
- throw new FileAlreadyExistsException(
- srcPath.toString(),
- dstPath.toString(),
- "Target object already exists and strategy is
FAIL_IF_EXISTS");
- default:
- throw new IllegalStateException("Unknown OverwriteStrategy: " +
dstOverwrite);
+ future.get();
+ } catch (CancellationException e) {
+ // Cancelled below after an earlier failure.
+ } catch (ExecutionException e) {
+ Throwable cause = e.getCause() != null ? e.getCause() : e;
+ if (failure == null) {
+ failure = cause;
+ for (Future<Void> other : futures) {
+ other.cancel(false);
+ }
+ } else {
+ failure.addSuppressed(cause);
}
}
}
-
+ if (failure != null) {
+ // Rethrow as is, so that callers still see the specific type (e.g.
FileNotFoundException,
+ // FileAlreadyExistsException, AccessDeniedException) the sequential
version threw.
+ if (failure instanceof IOException) {
+ throw (IOException) failure;
+ }
+ if (failure instanceof RuntimeException) {
+ throw (RuntimeException) failure;
+ }
+ if (failure instanceof Error) {
+ throw (Error) failure;
+ }
+ throw new IOException("Error rewriting objects", failure);
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new IOException("Interrupted while rewriting objects", e);
+ } finally {
+ // A cancelled future reports done even if its task is still running, so
wait for the pool
+ // itself: no file is left mid-rewrite when this method returns. If the
wait times out or the
+ // caller is interrupted (including above), stop waiting and interrupt
the in-flight rewrites.
+ executor.shutdown();
try {
- CopyWriter copyWriter = storage().copy(copyRequestBuilder.build());
- copyWriter.getResult();
-
- if (deleteSrc) {
- if (!storage().delete(srcId)) {
- // This may happen if the source file is deleted by another
process after copy.
- LOG.warn(
- "Source file {} could not be deleted after move to {}. It may
not have existed.",
- srcPath,
- dstPath);
- }
+ if (!executor.awaitTermination(REWRITE_TERMINATION_TIMEOUT_MINUTES,
TimeUnit.MINUTES)) {
+ LOG.warn(
+ "Timed out after {} minutes waiting for in-flight GCS rewrites
to finish; "
+ + "interrupting them. Some objects may be left partially
copied or moved.",
+ REWRITE_TERMINATION_TIMEOUT_MINUTES);
+ executor.shutdownNow();
}
+ } catch (InterruptedException e) {
+ executor.shutdownNow();
+ Thread.currentThread().interrupt();
+ }
+ }
+ }
+
+ /** Rewrites a single object, then deletes the source if {@code deleteSrc}.
*/
+ private void rewriteOne(
+ GcsPath srcPath,
+ GcsPath dstPath,
+ boolean deleteSrc,
+ MissingStrategy srcMissing,
+ OverwriteStrategy dstOverwrite)
+ throws IOException {
+ BlobId srcId = BlobId.of(srcPath.getBucket(), srcPath.getObject());
+ BlobId dstId = BlobId.of(dstPath.getBucket(), dstPath.getObject());
+
+ CopyRequest.Builder copyRequestBuilder =
+ CopyRequest.newBuilder()
+ .setSource(srcId)
+ .setMegabytesCopiedPerChunk(MEGABYTES_COPIED_PER_CHUNK);
+
+ if (dstOverwrite == OverwriteStrategy.ALWAYS_OVERWRITE) {
+ copyRequestBuilder.setTarget(dstId);
+ } else {
+ // FAIL_IF_EXISTS, SKIP_IF_EXISTS and SAFE_OVERWRITE require checking
the target blob
+ BlobInfo existingTarget;
+ try {
+ existingTarget = storage().get(dstId);
} catch (StorageException e) {
- if (e.getCode() == 404 && srcMissing ==
MissingStrategy.SKIP_IF_MISSING) {
+ throw translateStorageException(dstPath, e);
+ }
+
+ if (existingTarget == null) {
+ copyRequestBuilder.setTarget(dstId,
Storage.BlobTargetOption.doesNotExist());
+ } else {
+ switch (dstOverwrite) {
+ case SKIP_IF_EXISTS:
+ LOG.warn("Ignoring rewriting from {} to {} because target
exists.", srcPath, dstPath);
+ return; // Skip this file
+
+ case SAFE_OVERWRITE:
+ copyRequestBuilder.setTarget(
+ dstId,
Storage.BlobTargetOption.generationMatch(existingTarget.getGeneration()));
+ break;
+
+ case FAIL_IF_EXISTS:
+ throw new FileAlreadyExistsException(
+ srcPath.toString(),
+ dstPath.toString(),
+ "Target object already exists and strategy is FAIL_IF_EXISTS");
+ default:
+ throw new IllegalStateException("Unknown OverwriteStrategy: " +
dstOverwrite);
+ }
+ }
+ }
+
+ try {
+ CopyWriter copyWriter = storage().copy(copyRequestBuilder.build());
+ copyWriter.getResult();
+
+ if (deleteSrc) {
+ if (!storage().delete(srcId)) {
+ // This may happen if the source file is deleted by another process
after copy.
LOG.warn(
- "Ignoring rewriting from {} to {} because source does not
exist.", srcPath, dstPath);
- continue;
+ "Source file {} could not be deleted after move to {}. It may
not have existed.",
+ srcPath,
+ dstPath);
}
- throw translateStorageException(srcPath, e);
}
+ } catch (StorageException e) {
+ if (e.getCode() == 404 && srcMissing == MissingStrategy.SKIP_IF_MISSING)
{
+ LOG.warn(
+ "Ignoring rewriting from {} to {} because source does not exist.",
srcPath, dstPath);
+ return;
+ }
+ throw translateStorageException(srcPath, e);
}
}
diff --git
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
index b14a6fe602e..0654f974444 100644
---
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
+++
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
@@ -42,12 +42,15 @@ import com.google.cloud.WriteChannel;
import com.google.cloud.hadoop.util.AsyncWriteChannelOptions;
import com.google.cloud.http.HttpTransportOptions;
import com.google.cloud.storage.Blob;
+import com.google.cloud.storage.BlobId;
import com.google.cloud.storage.BlobInfo;
import com.google.cloud.storage.Bucket;
import com.google.cloud.storage.BucketInfo;
+import com.google.cloud.storage.CopyWriter;
import com.google.cloud.storage.Storage;
import com.google.cloud.storage.Storage.BlobWriteOption;
import com.google.cloud.storage.Storage.BucketGetOption;
+import com.google.cloud.storage.Storage.CopyRequest;
import com.google.cloud.storage.StorageBatch;
import com.google.cloud.storage.StorageBatchResult;
import com.google.cloud.storage.StorageException;
@@ -62,9 +65,12 @@ import java.nio.channels.SeekableByteChannel;
import java.nio.channels.WritableByteChannel;
import java.nio.file.AccessDeniedException;
import java.nio.file.FileAlreadyExistsException;
+import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
import
org.apache.beam.repackaged.core.org.apache.commons.compress.utils.SeekableInMemoryByteChannel;
import org.apache.beam.runners.core.metrics.CounterCell;
import org.apache.beam.runners.core.metrics.GcpResourceIdentifiers;
@@ -77,6 +83,7 @@ import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
import org.apache.beam.sdk.extensions.gcp.util.GcsUtil.CreateOptions;
import
org.apache.beam.sdk.extensions.gcp.util.GcsUtil.StorageObjectOrIOException;
import org.apache.beam.sdk.extensions.gcp.util.gcsfs.GcsPath;
+import org.apache.beam.sdk.io.fs.MoveOptions.StandardMoveOptions;
import org.apache.beam.sdk.metrics.MetricName;
import org.apache.beam.sdk.metrics.MetricsEnvironment;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
@@ -915,4 +922,72 @@ public class GcsUtilV2Test {
withForbidden.remove(
ImmutableList.of(deleted.toString(), missing.toString(),
forbidden.toString())));
}
+
+ /**
+ * java-storage cannot batch rewrites, so V2 issues the files of a rename
concurrently. Every
+ * rewrite here blocks until all of them have started, which only happens if
they overlap.
+ */
+ @Test
+ public void testV2RenameRewritesFilesConcurrently() throws IOException {
+ int numFiles = 4;
+ CountDownLatch allStarted = new CountDownLatch(numFiles);
+ Storage storage = Mockito.mock(Storage.class);
+ when(storage.copy(any(CopyRequest.class)))
+ .thenAnswer(
+ invocation -> {
+ allStarted.countDown();
+ assertTrue("rewrites did not overlap", allStarted.await(30,
TimeUnit.SECONDS));
+ return Mockito.mock(CopyWriter.class);
+ });
+ when(storage.delete(any(BlobId.class))).thenReturn(true);
+ GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+ List<String> srcs = new ArrayList<>();
+ List<String> dsts = new ArrayList<>();
+ for (int i = 0; i < numFiles; i++) {
+ srcs.add("gs://testbucket/src" + i);
+ dsts.add("gs://testbucket/dst" + i);
+ }
+
+ gcsUtil.rename(srcs, dsts);
+
+ verify(storage, Mockito.times(numFiles)).copy(any(CopyRequest.class));
+ for (int i = 0; i < numFiles; i++) {
+ verify(storage).delete(BlobId.of("testbucket", "src" + i));
+ }
+ }
+
+ /**
+ * Running the files concurrently must not change what a rename reports: a
missing source is
+ * skipped or fails depending on IGNORE_MISSING_FILES, and other failures
keep their type.
+ */
+ @Test
+ public void testV2RenameFailuresAreReportedPerFile() throws IOException {
+ Storage storage = Mockito.mock(Storage.class);
+ when(storage.copy(any(CopyRequest.class)))
+ .thenAnswer(
+ invocation -> {
+ String source =
invocation.<CopyRequest>getArgument(0).getSource().getName();
+ if (source.equals("missing")) {
+ throw new StorageException(404, "Not Found");
+ }
+ if (source.equals("forbidden")) {
+ throw new StorageException(403, "Forbidden");
+ }
+ return Mockito.mock(CopyWriter.class);
+ });
+ GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+ List<String> okAndMissing = ImmutableList.of("gs://testbucket/ok",
"gs://testbucket/missing");
+ List<String> okAndMissingDsts =
+ ImmutableList.of("gs://testbucket/ok.bak",
"gs://testbucket/missing.bak");
+
+ gcsUtil.rename(okAndMissing, okAndMissingDsts,
StandardMoveOptions.IGNORE_MISSING_FILES);
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.rename(okAndMissing, okAndMissingDsts));
+ assertThrows(
+ AccessDeniedException.class,
+ () ->
+ gcsUtil.rename(
+ ImmutableList.of("gs://testbucket/ok",
"gs://testbucket/forbidden"),
+ ImmutableList.of("gs://testbucket/ok.bak",
"gs://testbucket/forbidden.bak"),
+ StandardMoveOptions.IGNORE_MISSING_FILES));
+ }
}