Skip to content

Commit 60ee96f

Browse files
committed
[CELEBORN-2174] Remove partitionSplitEnabled of ReserveSlots and FileInfo
1 parent 9482808 commit 60ee96f

25 files changed

Lines changed: 33 additions & 138 deletions

File tree

‎client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1302,7 +1302,6 @@ class LifecycleManager(val appUniqueId: String, val conf: CelebornConf) extends
13021302
rangeReadFilter,
13031303
userIdentifier,
13041304
conf.pushDataTimeoutMs,
1305-
partitionSplitEnabled = true,
13061305
isSegmentGranularityVisible = isSegmentGranularityVisible))
13071306
futures.add((future, workerInfo))
13081307
}(ec)

‎common/src/main/java/org/apache/celeborn/common/meta/DiskFileInfo.java‎

Lines changed: 6 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@
1919

2020
import java.io.File;
2121
import java.util.ArrayList;
22-
import java.util.Arrays;
22+
import java.util.Collections;
2323

2424
import com.google.common.annotations.VisibleForTesting;
2525
import org.apache.hadoop.fs.FileSystem;
@@ -42,23 +42,18 @@ public class DiskFileInfo extends FileInfo {
4242

4343
public DiskFileInfo(
4444
UserIdentifier userIdentifier,
45-
boolean partitionSplitEnabled,
4645
FileMeta fileMeta,
4746
String filePath,
4847
StorageInfo.Type storageType) {
49-
super(userIdentifier, partitionSplitEnabled, fileMeta);
48+
super(userIdentifier, fileMeta);
5049
this.filePath = filePath;
5150
this.storageType = storageType;
5251
}
5352

5453
// only called when restore from pb or in UT
5554
public DiskFileInfo(
56-
UserIdentifier userIdentifier,
57-
boolean partitionSplitEnabled,
58-
FileMeta fileMeta,
59-
String filePath,
60-
long bytesFlushed) {
61-
super(userIdentifier, partitionSplitEnabled, fileMeta);
55+
UserIdentifier userIdentifier, FileMeta fileMeta, String filePath, long bytesFlushed) {
56+
super(userIdentifier, fileMeta);
6257
this.filePath = filePath;
6358
this.storageType = StorageInfo.Type.HDD;
6459
this.bytesFlushed = bytesFlushed;
@@ -68,14 +63,13 @@ public DiskFileInfo(
6863
public DiskFileInfo(File file, UserIdentifier userIdentifier, CelebornConf conf) {
6964
this(
7065
userIdentifier,
71-
true,
72-
new ReduceFileMeta(new ArrayList<>(Arrays.asList(0L)), conf.shuffleChunkSize()),
66+
new ReduceFileMeta(new ArrayList<>(Collections.singletonList(0L)), conf.shuffleChunkSize()),
7367
file.getAbsolutePath(),
7468
StorageInfo.Type.HDD);
7569
}
7670

7771
public DiskFileInfo(UserIdentifier userIdentifier, FileMeta fileMeta, String filePath) {
78-
super(userIdentifier, true, fileMeta);
72+
super(userIdentifier, fileMeta);
7973
this.filePath = filePath;
8074
this.storageType = StorageInfo.Type.HDD;
8175
}

‎common/src/main/java/org/apache/celeborn/common/meta/FileInfo.java‎

Lines changed: 1 addition & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -24,17 +24,13 @@
2424

2525
public abstract class FileInfo {
2626
private final UserIdentifier userIdentifier;
27-
// whether to split is decided by client side.
28-
// now it's just used for mappartition to compatible with old client which can't support split
29-
private boolean partitionSplitEnabled;
3027
protected FileMeta fileMeta;
3128
protected final Set<Long> streams = ConcurrentHashMap.newKeySet();
3229
protected volatile long bytesFlushed;
3330
private boolean isReduceFileMeta;
3431

35-
public FileInfo(UserIdentifier userIdentifier, boolean partitionSplitEnabled, FileMeta fileMeta) {
32+
public FileInfo(UserIdentifier userIdentifier, FileMeta fileMeta) {
3633
this.userIdentifier = userIdentifier;
37-
this.partitionSplitEnabled = partitionSplitEnabled;
3834
this.fileMeta = fileMeta;
3935
this.isReduceFileMeta = fileMeta instanceof ReduceFileMeta;
4036
}
@@ -67,14 +63,6 @@ public UserIdentifier getUserIdentifier() {
6763
return userIdentifier;
6864
}
6965

70-
public boolean isPartitionSplitEnabled() {
71-
return partitionSplitEnabled;
72-
}
73-
74-
public void setPartitionSplitEnabled(boolean partitionSplitEnabled) {
75-
this.partitionSplitEnabled = partitionSplitEnabled;
76-
}
77-
7866
public boolean addStream(long streamId) {
7967
if (!isReduceFileMeta) {
8068
throw new IllegalStateException("In addStream, filemeta cannot be MapFileMeta");

‎common/src/main/java/org/apache/celeborn/common/meta/MemoryFileInfo.java‎

Lines changed: 4 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -32,18 +32,13 @@ public class MemoryFileInfo extends FileInfo {
3232
private CompositeByteBuf sortedBuffer;
3333
private Map<Integer, List<ShuffleBlockInfo>> sortedIndexes;
3434

35-
public MemoryFileInfo(
36-
UserIdentifier userIdentifier, boolean partitionSplitEnabled, FileMeta fileMeta) {
37-
super(userIdentifier, partitionSplitEnabled, fileMeta);
35+
public MemoryFileInfo(UserIdentifier userIdentifier, FileMeta fileMeta) {
36+
super(userIdentifier, fileMeta);
3837
}
3938

4039
// This constructor is only used in partition sorter for temp new memory file
41-
public MemoryFileInfo(
42-
UserIdentifier userIdentifier,
43-
boolean partitionSplitEnabled,
44-
FileMeta fileMeta,
45-
CompositeByteBuf buffer) {
46-
super(userIdentifier, partitionSplitEnabled, fileMeta);
40+
public MemoryFileInfo(UserIdentifier userIdentifier, FileMeta fileMeta, CompositeByteBuf buffer) {
41+
super(userIdentifier, fileMeta);
4742
this.buffer = buffer;
4843
}
4944

‎common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala‎

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -526,7 +526,6 @@ object ControlMessages extends Logging {
526526
rangeReadFilter: Boolean,
527527
userIdentifier: UserIdentifier,
528528
pushDataTimeout: Long,
529-
partitionSplitEnabled: Boolean = false,
530529
isSegmentGranularityVisible: Boolean = false)
531530
extends WorkerMessage
532531

@@ -967,7 +966,6 @@ object ControlMessages extends Logging {
967966
rangeReadFilter,
968967
userIdentifier,
969968
pushDataTimeout,
970-
partitionSplitEnabled,
971969
isSegmentGranularityVisible) =>
972970
val payload = PbReserveSlots.newBuilder()
973971
.setApplicationId(applicationId)
@@ -980,7 +978,6 @@ object ControlMessages extends Logging {
980978
.setRangeReadFilter(rangeReadFilter)
981979
.setUserIdentifier(PbSerDeUtils.toPbUserIdentifier(userIdentifier))
982980
.setPushDataTimeout(pushDataTimeout)
983-
.setPartitionSplitEnabled(partitionSplitEnabled)
984981
.setIsSegmentGranularityVisible(isSegmentGranularityVisible)
985982
.build().toByteArray
986983
new TransportMessage(MessageType.RESERVE_SLOTS, payload)
@@ -1391,7 +1388,6 @@ object ControlMessages extends Logging {
13911388
pbReserveSlots.getRangeReadFilter,
13921389
userIdentifier,
13931390
pbReserveSlots.getPushDataTimeout,
1394-
pbReserveSlots.getPartitionSplitEnabled,
13951391
pbReserveSlots.getIsSegmentGranularityVisible)
13961392

13971393
case RESERVE_SLOTS_RESPONSE_VALUE =>

‎common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -116,7 +116,6 @@ object PbSerDeUtils {
116116
}
117117
new DiskFileInfo(
118118
userIdentifier,
119-
pbFileInfo.getPartitionSplitEnabled,
120119
meta,
121120
pbFileInfo.getFilePath,
122121
pbFileInfo.getBytesFlushed)
@@ -140,7 +139,6 @@ object PbSerDeUtils {
140139
.setFilePath(fileInfo.getFilePath)
141140
.setUserIdentifier(toPbUserIdentifier(fileInfo.getUserIdentifier))
142141
.setBytesFlushed(fileInfo.getFileLength)
143-
.setPartitionSplitEnabled(fileInfo.isPartitionSplitEnabled)
144142
if (fileInfo.getFileMeta.isInstanceOf[MapFileMeta]) {
145143
val mapFileMeta = fileInfo.getFileMeta.asInstanceOf[MapFileMeta]
146144
builder.setPartitionType(PartitionType.MAP.getValue)

‎common/src/test/scala/org/apache/celeborn/common/util/PbSerDeUtilsTest.scala‎

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -71,25 +71,21 @@ class PbSerDeUtilsTest extends CelebornFunSuite {
7171

7272
val fileInfo1 = new DiskFileInfo(
7373
userIdentifier1,
74-
true,
7574
new ReduceFileMeta(chunkOffsets1, 123),
7675
file1.getAbsolutePath,
7776
3000L)
7877
val fileInfo2 = new DiskFileInfo(
7978
userIdentifier2,
80-
true,
8179
new ReduceFileMeta(chunkOffsets2, 123),
8280
file2.getAbsolutePath,
8381
6000L)
8482
val mapFileInfo1 = new DiskFileInfo(
8583
userIdentifier1,
86-
true,
8784
new MapFileMeta(1024, 10),
8885
file1.getAbsolutePath,
8986
6000L)
9087
val mapFileInfo2 = new DiskFileInfo(
9188
userIdentifier2,
92-
true,
9389
new MapFileMeta(1024, 10),
9490
file2.getAbsolutePath,
9591
6000L)

‎worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/MapPartitionDataReader.java‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -455,7 +455,7 @@ public void release() {
455455
logger.debug("release reader for stream {}", streamId);
456456
// old client can't support BufferStreamEnd, so for new client it tells client that this
457457
// stream is finished.
458-
if (fileInfo.isPartitionSplitEnabled() && !errorNotified) {
458+
if (!errorNotified) {
459459
associatedChannel.writeAndFlush(
460460
new RpcRequest(
461461
TransportClient.requestId(),

‎worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionDataWriterContext.java‎

Lines changed: 0 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,6 @@ public class PartitionDataWriterContext {
3333
private final String appId;
3434
private final int shuffleId;
3535
private final UserIdentifier userIdentifier;
36-
private final boolean partitionSplitEnabled;
3736
private final String shuffleKey;
3837
private final PartitionType partitionType;
3938
private final boolean isSegmentGranularityVisible;
@@ -51,7 +50,6 @@ public PartitionDataWriterContext(
5150
int shuffleId,
5251
UserIdentifier userIdentifier,
5352
PartitionType partitionType,
54-
boolean partitionSplitEnabled,
5553
boolean isSegmentGranularityVisible) {
5654
this.splitThreshold = splitThreshold;
5755
this.partitionSplitMode = partitionSplitMode;
@@ -60,7 +58,6 @@ public PartitionDataWriterContext(
6058
this.appId = appId;
6159
this.shuffleId = shuffleId;
6260
this.userIdentifier = userIdentifier;
63-
this.partitionSplitEnabled = partitionSplitEnabled;
6461
this.partitionType = partitionType;
6562
this.shuffleKey = Utils.makeShuffleKey(appId, shuffleId);
6663
this.isSegmentGranularityVisible = isSegmentGranularityVisible;
@@ -94,10 +91,6 @@ public UserIdentifier getUserIdentifier() {
9491
return userIdentifier;
9592
}
9693

97-
public boolean isPartitionSplitEnabled() {
98-
return partitionSplitEnabled;
99-
}
100-
10194
public String getShuffleKey() {
10295
return shuffleKey;
10396
}
@@ -152,8 +145,6 @@ public String toString() {
152145
+ shuffleId
153146
+ ", userIdentifier="
154147
+ userIdentifier
155-
+ ", partitionSplitEnabled="
156-
+ partitionSplitEnabled
157148
+ ", shuffleKey='"
158149
+ shuffleKey
159150
+ '\''

‎worker/src/main/java/org/apache/celeborn/service/deploy/worker/storage/PartitionFilesSorter.java‎

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -227,11 +227,7 @@ public FileInfo getSortedFileInfo(
227227
memoryFileInfo.getSortedBuffer(),
228228
targetBuffer,
229229
shuffleChunkSize);
230-
return new MemoryFileInfo(
231-
memoryFileInfo.getUserIdentifier(),
232-
memoryFileInfo.isPartitionSplitEnabled(),
233-
reduceFileMeta,
234-
targetBuffer);
230+
return new MemoryFileInfo(memoryFileInfo.getUserIdentifier(), reduceFileMeta, targetBuffer);
235231
} else {
236232
DiskFileInfo diskFileInfo = ((DiskFileInfo) fileInfo);
237233
String fileId = shuffleKey + "-" + fileName;

0 commit comments

Comments
 (0)