This is an automated email from the ASF dual-hosted git repository.

angerszhuuuu pushed a commit to branch branch-0.4
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git


The following commit(s) were added to refs/heads/branch-0.4 by this push:
     new 0038b0f05 [CELEBORN-1226][FOLLOWUP] Unify creation of thread using 
ThreadUtils
0038b0f05 is described below

commit 0038b0f057fbf58809e38ac36017b8b2d7a16a87
Author: Angerszhuuuu <[email protected]>
AuthorDate: Tue Jan 23 10:39:58 2024 +0800

    [CELEBORN-1226][FOLLOWUP] Unify creation of thread using ThreadUtils
    
    ### What changes were proposed in this pull request?
    Unify creation of thread using ThreadUtils
    
    ### Why are the changes needed?
    
    ### Does this PR introduce _any_ user-facing change?
    
    ### How was this patch tested?
    
    Closes #2247 from AngersZhuuuu/CELEBORN-1226-FOLLOWUP.
    
    Authored-by: Angerszhuuuu <[email protected]>
    Signed-off-by: Angerszhuuuu <[email protected]>
    (cherry picked from commit cee251e683694a445e7ba3f13db78d8ed8666870)
    Signed-off-by: Angerszhuuuu <[email protected]>
---
 .../celeborn/client/read/DfsPartitionReader.java   | 104 +++++++++++----------
 1 file changed, 53 insertions(+), 51 deletions(-)

diff --git 
a/client/src/main/java/org/apache/celeborn/client/read/DfsPartitionReader.java 
b/client/src/main/java/org/apache/celeborn/client/read/DfsPartitionReader.java
index d379a1139..458e5d0e5 100644
--- 
a/client/src/main/java/org/apache/celeborn/client/read/DfsPartitionReader.java
+++ 
b/client/src/main/java/org/apache/celeborn/client/read/DfsPartitionReader.java
@@ -21,6 +21,7 @@ import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.List;
+import java.util.concurrent.ExecutorService;
 import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
@@ -45,23 +46,25 @@ import org.apache.celeborn.common.protocol.PbOpenStream;
 import org.apache.celeborn.common.protocol.PbStreamHandler;
 import org.apache.celeborn.common.protocol.StreamType;
 import org.apache.celeborn.common.util.ShuffleBlockInfoUtils;
-import org.apache.celeborn.common.util.ThreadExceptionHandler;
+import org.apache.celeborn.common.util.ThreadUtils;
 import org.apache.celeborn.common.util.Utils;
 
 public class DfsPartitionReader implements PartitionReader {
   private static Logger logger = 
LoggerFactory.getLogger(DfsPartitionReader.class);
+  private CelebornConf conf;
   PartitionLocation location;
   private final long shuffleChunkSize;
   private final int fetchMaxReqsInFlight;
   private final LinkedBlockingQueue<ByteBuf> results;
   private final AtomicReference<IOException> exception = new 
AtomicReference<>();
   private volatile boolean closed = false;
-  private Thread fetchThread;
+  private ExecutorService fetchThread;
   private boolean fetchThreadStarted;
   private FSDataInputStream hdfsInputStream;
   private int numChunks = 0;
   private int returnedChunks = 0;
   private int currentChunkIndex = 0;
+  private final List<Long> chunkOffsets = new ArrayList<>();
   private TransportClient client;
   private PbStreamHandler streamHandler;
   private MetricsCallback metricsCallback;
@@ -75,6 +78,7 @@ public class DfsPartitionReader implements PartitionReader {
       int endMapIndex,
       MetricsCallback metricsCallback)
       throws IOException {
+    this.conf = conf;
     shuffleChunkSize = conf.dfsReadChunkSize();
     fetchMaxReqsInFlight = conf.clientFetchMaxReqsInFlight();
     results = new LinkedBlockingQueue<>();
@@ -82,7 +86,6 @@ public class DfsPartitionReader implements PartitionReader {
     this.metricsCallback = metricsCallback;
     this.location = location;
 
-    final List<Long> chunkOffsets = new ArrayList<>();
     if (endMapIndex != Integer.MAX_VALUE) {
       long fetchTimeoutMs = conf.clientFetchTimeoutMs();
       try {
@@ -124,53 +127,8 @@ public class DfsPartitionReader implements PartitionReader 
{
     if (chunkOffsets.size() > 1) {
       numChunks = chunkOffsets.size() - 1;
       fetchThread =
-          new Thread(
-              () -> {
-                try {
-                  while (!closed && currentChunkIndex < numChunks) {
-                    while (results.size() >= fetchMaxReqsInFlight) {
-                      Thread.sleep(50);
-                    }
-                    long offset = chunkOffsets.get(currentChunkIndex);
-                    long length = chunkOffsets.get(currentChunkIndex + 1) - 
offset;
-                    logger.debug("read {} offset {} length {}", 
currentChunkIndex, offset, length);
-                    byte[] buffer = new byte[(int) length];
-                    try {
-                      hdfsInputStream.readFully(offset, buffer);
-                    } catch (IOException e) {
-                      logger.warn(
-                          "read HDFS {} failed will retry, error detail {}",
-                          location.getStorageInfo().getFilePath(),
-                          e);
-                      try {
-                        hdfsInputStream.close();
-                        hdfsInputStream =
-                            ShuffleClient.getHdfsFs(conf)
-                                .open(
-                                    new Path(
-                                        Utils.getSortedFilePath(
-                                            
location.getStorageInfo().getFilePath())));
-                        hdfsInputStream.readFully(offset, buffer);
-                      } catch (IOException ex) {
-                        logger.warn(
-                            "retry read HDFS {} failed, error detail {} ",
-                            location.getStorageInfo().getFilePath(),
-                            e);
-                        exception.set(ex);
-                        break;
-                      }
-                    }
-                    results.put(Unpooled.wrappedBuffer(buffer));
-                    logger.debug("add index {} to results", 
currentChunkIndex++);
-                  }
-                } catch (Exception e) {
-                  logger.warn("Fetch thread is cancelled.", e);
-                  // cancel a task for speculative, ignore this exception
-                }
-                logger.debug("fetch {} is done.", 
location.getStorageInfo().getFilePath());
-              },
+          ThreadUtils.newDaemonSingleThreadExecutor(
               "Dfs-fetch-thread" + location.getStorageInfo().getFilePath());
-      fetchThread.setUncaughtExceptionHandler(new 
ThreadExceptionHandler(fetchThread.getName()));
       logger.debug("Start dfs read on location {}", location);
       ShuffleClient.incrementTotalReadCounter();
     }
@@ -225,7 +183,51 @@ public class DfsPartitionReader implements PartitionReader 
{
     ByteBuf chunk = null;
     if (!fetchThreadStarted) {
       fetchThreadStarted = true;
-      fetchThread.start();
+      fetchThread.submit(
+          () -> {
+            try {
+              while (!closed && currentChunkIndex < numChunks) {
+                while (results.size() >= fetchMaxReqsInFlight) {
+                  Thread.sleep(50);
+                }
+                long offset = chunkOffsets.get(currentChunkIndex);
+                long length = chunkOffsets.get(currentChunkIndex + 1) - offset;
+                logger.debug("read {} offset {} length {}", currentChunkIndex, 
offset, length);
+                byte[] buffer = new byte[(int) length];
+                try {
+                  hdfsInputStream.readFully(offset, buffer);
+                } catch (IOException e) {
+                  logger.warn(
+                      "read HDFS {} failed will retry, error detail {}",
+                      location.getStorageInfo().getFilePath(),
+                      e);
+                  try {
+                    hdfsInputStream.close();
+                    hdfsInputStream =
+                        ShuffleClient.getHdfsFs(conf)
+                            .open(
+                                new Path(
+                                    Utils.getSortedFilePath(
+                                        
location.getStorageInfo().getFilePath())));
+                    hdfsInputStream.readFully(offset, buffer);
+                  } catch (IOException ex) {
+                    logger.warn(
+                        "retry read HDFS {} failed, error detail {} ",
+                        location.getStorageInfo().getFilePath(),
+                        e);
+                    exception.set(ex);
+                    break;
+                  }
+                }
+                results.put(Unpooled.wrappedBuffer(buffer));
+                logger.debug("add index {} to results", currentChunkIndex++);
+              }
+            } catch (Exception e) {
+              logger.warn("Fetch thread is cancelled.", e);
+              // cancel a task for speculative, ignore this exception
+            }
+            logger.debug("fetch {} is done.", 
location.getStorageInfo().getFilePath());
+          });
     }
     try {
       while (chunk == null) {
@@ -255,7 +257,7 @@ public class DfsPartitionReader implements PartitionReader {
   public void close() {
     closed = true;
     if (fetchThread != null) {
-      fetchThread.interrupt();
+      fetchThread.shutdownNow();
     }
     try {
       hdfsInputStream.close();

Reply via email to