shunping commented on code in PR #40376:
URL: https://github.com/apache/beam/pull/40376#discussion_r4209624321
##########
sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java:
##########
@@ -594,73 +604,159 @@ private void rewriteHelper(
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(), MAX_CONCURRENT_REWRITES);
+ 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.
+ @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);
Review Comment:
Revised the code to better handling shutdown on executor. Waiting to push
the change but github remote server is not healthy right now.
--
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]