SteNicholas commented on code in PR #3784:
URL: https://github.com/apache/celeborn/pull/3784#discussion_r3845601589


##########
worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java:
##########
@@ -314,6 +327,174 @@ public FileInfo getSortedFileInfo(
     }
   }
 
+  public CompletableFuture<FileInfo> getSortedFileInfoAsync(
+      String shuffleKey, String fileName, FileInfo fileInfo, int 
startMapIndex, int endMapIndex) {
+    if (shutdown) {
+      return failedSortedFileInfo(new IOException("Partition sorter is 
closed."));
+    }
+    if (fileInfo instanceof MemoryFileInfo) {
+      logger.debug(
+          "Sorting memory shuffle file synchronously for shuffle key {}, file 
name {}, map range [{}, {})",
+          shuffleKey,
+          fileName,
+          startMapIndex,
+          endMapIndex);
+      try {
+        return CompletableFuture.completedFuture(
+            getSortedFileInfo(shuffleKey, fileName, fileInfo, startMapIndex, 
endMapIndex));
+      } catch (IOException e) {
+        return failedSortedFileInfo(e);
+      }
+    }
+
+    DiskFileInfo diskFileInfo = (DiskFileInfo) fileInfo;
+    String fileId = shuffleKey + "-" + fileName;
+    UserIdentifier userIdentifier = diskFileInfo.getUserIdentifier();
+    Set<String> sorted =
+        sortedShuffleFiles.computeIfAbsent(shuffleKey, ignored -> 
ConcurrentHashMap.newKeySet());
+    Set<String> sorting =
+        sortingShuffleFiles.computeIfAbsent(shuffleKey, ignored -> 
ConcurrentHashMap.newKeySet());
+    String sortedFilePath = 
Utils.getSortedFilePath(diskFileInfo.getFilePath());
+    String indexFilePath = Utils.getIndexFilePath(diskFileInfo.getFilePath());
+
+    CompletableFuture<Void> sortCompletionFuture;
+    synchronized (sorting) {
+      if (shutdown) {
+        return failedSortedFileInfo(new IOException("Partition sorter is 
closed."));
+      }
+      if (sorted.contains(fileId)) {
+        sortCompletionFuture = CompletableFuture.completedFuture(null);
+      } else {
+        try {
+          sortCompletionFuture =
+              sortCompletionFutures
+                  .computeIfAbsent(shuffleKey, ignored -> 
JavaUtils.newConcurrentHashMap())
+                  .compute(
+                      fileId,
+                      (ignored, existing) ->
+                          existing == null || existing.isDone()
+                              ? createSortCompletionFuture(shuffleKey, fileId, 
diskFileInfo)

Review Comment:
   **[Major] Preserve or explicitly define the timeout semantics.** All waiters 
reuse this future and therefore share one timeout that starts with the first 
waiter. A later request inherits only the remaining time, potentially almost 
zero, while a request arriving just after this future times out gets a fresh 
full timeout. The previous synchronous path applied the timeout per request. 
Please share the raw sort-completion future but apply a timeout to each 
returned waiter, or document and test the new global sort-attempt timeout 
semantics.



##########
worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java:
##########
@@ -314,6 +327,174 @@ public FileInfo getSortedFileInfo(
     }
   }
 
+  public CompletableFuture<FileInfo> getSortedFileInfoAsync(
+      String shuffleKey, String fileName, FileInfo fileInfo, int 
startMapIndex, int endMapIndex) {
+    if (shutdown) {
+      return failedSortedFileInfo(new IOException("Partition sorter is 
closed."));
+    }
+    if (fileInfo instanceof MemoryFileInfo) {
+      logger.debug(
+          "Sorting memory shuffle file synchronously for shuffle key {}, file 
name {}, map range [{}, {})",
+          shuffleKey,
+          fileName,
+          startMapIndex,
+          endMapIndex);
+      try {
+        return CompletableFuture.completedFuture(
+            getSortedFileInfo(shuffleKey, fileName, fileInfo, startMapIndex, 
endMapIndex));
+      } catch (IOException e) {
+        return failedSortedFileInfo(e);
+      }
+    }
+
+    DiskFileInfo diskFileInfo = (DiskFileInfo) fileInfo;
+    String fileId = shuffleKey + "-" + fileName;
+    UserIdentifier userIdentifier = diskFileInfo.getUserIdentifier();
+    Set<String> sorted =
+        sortedShuffleFiles.computeIfAbsent(shuffleKey, ignored -> 
ConcurrentHashMap.newKeySet());
+    Set<String> sorting =
+        sortingShuffleFiles.computeIfAbsent(shuffleKey, ignored -> 
ConcurrentHashMap.newKeySet());
+    String sortedFilePath = 
Utils.getSortedFilePath(diskFileInfo.getFilePath());
+    String indexFilePath = Utils.getIndexFilePath(diskFileInfo.getFilePath());
+
+    CompletableFuture<Void> sortCompletionFuture;
+    synchronized (sorting) {
+      if (shutdown) {
+        return failedSortedFileInfo(new IOException("Partition sorter is 
closed."));
+      }
+      if (sorted.contains(fileId)) {
+        sortCompletionFuture = CompletableFuture.completedFuture(null);
+      } else {
+        try {
+          sortCompletionFuture =
+              sortCompletionFutures
+                  .computeIfAbsent(shuffleKey, ignored -> 
JavaUtils.newConcurrentHashMap())
+                  .compute(
+                      fileId,
+                      (ignored, existing) ->
+                          existing == null || existing.isDone()
+                              ? createSortCompletionFuture(shuffleKey, fileId, 
diskFileInfo)
+                              : existing);
+        } catch (RejectedExecutionException e) {
+          return failedSortedFileInfo(new IOException("Partition sorter is 
closed.", e));
+        }
+        if (!sorting.contains(fileId)) {
+          try {
+            FileSorter fileSorter = new FileSorter(diskFileInfo, fileId, 
shuffleKey);

Review Comment:
   **[Major] Move `FileSorter` construction off the fetch thread.** This method 
is called by the asynchronous open-stream path, but the constructor 
synchronously performs local `exists` and `delete` calls and, for DFS storage, 
`hadoopFs.exists` and `hadoopFs.delete`. Remote metadata operations can 
therefore block the fetch event loop before work reaches the sorter executor, 
undermining the goal of this PR. Please move construction and file-system 
preparation onto a worker executor.



##########
worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java:
##########
@@ -314,6 +327,174 @@ public FileInfo getSortedFileInfo(
     }
   }
 
+  public CompletableFuture<FileInfo> getSortedFileInfoAsync(
+      String shuffleKey, String fileName, FileInfo fileInfo, int 
startMapIndex, int endMapIndex) {
+    if (shutdown) {
+      return failedSortedFileInfo(new IOException("Partition sorter is 
closed."));
+    }
+    if (fileInfo instanceof MemoryFileInfo) {
+      logger.debug(
+          "Sorting memory shuffle file synchronously for shuffle key {}, file 
name {}, map range [{}, {})",
+          shuffleKey,
+          fileName,
+          startMapIndex,
+          endMapIndex);
+      try {
+        return CompletableFuture.completedFuture(
+            getSortedFileInfo(shuffleKey, fileName, fileInfo, startMapIndex, 
endMapIndex));
+      } catch (IOException e) {
+        return failedSortedFileInfo(e);
+      }
+    }
+
+    DiskFileInfo diskFileInfo = (DiskFileInfo) fileInfo;
+    String fileId = shuffleKey + "-" + fileName;
+    UserIdentifier userIdentifier = diskFileInfo.getUserIdentifier();
+    Set<String> sorted =
+        sortedShuffleFiles.computeIfAbsent(shuffleKey, ignored -> 
ConcurrentHashMap.newKeySet());
+    Set<String> sorting =
+        sortingShuffleFiles.computeIfAbsent(shuffleKey, ignored -> 
ConcurrentHashMap.newKeySet());
+    String sortedFilePath = 
Utils.getSortedFilePath(diskFileInfo.getFilePath());
+    String indexFilePath = Utils.getIndexFilePath(diskFileInfo.getFilePath());
+
+    CompletableFuture<Void> sortCompletionFuture;
+    synchronized (sorting) {
+      if (shutdown) {
+        return failedSortedFileInfo(new IOException("Partition sorter is 
closed."));
+      }
+      if (sorted.contains(fileId)) {
+        sortCompletionFuture = CompletableFuture.completedFuture(null);
+      } else {
+        try {
+          sortCompletionFuture =
+              sortCompletionFutures
+                  .computeIfAbsent(shuffleKey, ignored -> 
JavaUtils.newConcurrentHashMap())
+                  .compute(
+                      fileId,
+                      (ignored, existing) ->
+                          existing == null || existing.isDone()
+                              ? createSortCompletionFuture(shuffleKey, fileId, 
diskFileInfo)
+                              : existing);
+        } catch (RejectedExecutionException e) {
+          return failedSortedFileInfo(new IOException("Partition sorter is 
closed.", e));
+        }
+        if (!sorting.contains(fileId)) {
+          try {
+            FileSorter fileSorter = new FileSorter(diskFileInfo, fileId, 
shuffleKey);
+            sorting.add(fileId);
+            logger.debug(
+                "Adding sorter to sort queue shuffle key {}, file name {}", 
shuffleKey, fileName);
+            shuffleSortTaskDeque.put(fileSorter);
+          } catch (InterruptedException e) {
+            logger.error(
+                "Sorter scheduler thread is interrupted means worker is 
shutting down.", e);
+            sorting.remove(fileId);
+            sortCompletionFuture.completeExceptionally(
+                new IOException(
+                    "Sort scheduler thread is interrupted means worker is 
shutting down.", e));
+          } catch (IOException e) {
+            logger.error("File sorter access DFS failed.", e);
+            sortCompletionFuture.completeExceptionally(
+                new IOException("File sorter access DFS failed.", e));
+          }
+        }
+      }
+    }
+
+    if (shutdown) {
+      sortCompletionFuture.completeExceptionally(new IOException("Partition 
sorter is closed."));
+    }
+    sortedFileWaiterCount.incrementAndGet();
+    return sortCompletionFuture

Review Comment:
   **[Major] Track the resolve stage through shuffle cleanup.** 
`sortCompletionFuture` is removed from `sortCompletionFutures` as soon as 
sorting completes, while the future returned by `thenApplyAsync(resolve)` is 
not tracked. If shuffle cleanup runs after sorting but while `resolve` is 
queued or running, `cleanup` cannot fail this waiter; the continuation can 
still succeed, repopulate the index cache, and let `FetchHandler` register a 
stream after its one-time stream cleanup. I reproduced this by blocking 
`resolve`, invoking cleanup, and observing the waiter succeed after release. 
Please track or cancel the derived request future, or validate a per-shuffle 
generation before returning and registering the result.



-- 
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]

Reply via email to