This is an automated email from the ASF dual-hosted git repository.
angerszhuuuu pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git
The following commit(s) were added to refs/heads/main by this push:
new cee251e68 [CELEBORN-1226][FOLLOWUP] Unify creation of thread using
ThreadUtils
cee251e68 is described below
commit cee251e683694a445e7ba3f13db78d8ed8666870
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]>
---
.../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();