This is an automated email from the ASF dual-hosted git repository.
SteNicholas pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/celeborn.git
The following commit(s) were added to refs/heads/main by this push:
new 749e45522 [CELEBORN-2060][FOLLOWUP] Fix rack-aware replica storage
selection
749e45522 is described below
commit 749e45522a1726d25f5949d6ae82f1e30ac9f47e
Author: Kalvin2077 <[email protected]>
AuthorDate: Tue Aug 25 00:07:24 2026 +0800
[CELEBORN-2060][FOLLOWUP] Fix rack-aware replica storage selection
### What changes were proposed in this pull request?
This follow-up to [PR #3781](https://github.com/apache/celeborn/pull/3781):
- Validates that rack-aware replica candidates support the requested
storage type.
- Builds replica `StorageInfo` from the selected replica worker.
- Adds regression tests for diskless replicas and different mount points.
### Why are the changes needed?
During best-effort fallback, the allocator could select a diskless worker
for a local-disk replica or incorrectly reuse the primary worker's mount point.
### Does this PR resolve a correctness bug?
- [x] Yes
### Does this PR introduce _any_ user-facing change?
- [ ] Yes
### How was this patch tested?
`build/mvn -pl master -am -Dtest=SlotsAllocatorRackAwareSuiteJ test`
Closes #3803 from Kalvin2077/fix/CELEBORN-2423.
Authored-by: Kalvin2077 <[email protected]>
Signed-off-by: Nicholas Jiang <[email protected]>
---
.../deploy/master/slotsalloc/SlotsAllocator.java | 6 +-
.../slotsalloc/SlotsAllocatorRackAwareSuiteJ.java | 89 ++++++++++++++++++++++
2 files changed, 91 insertions(+), 4 deletions(-)
diff --git
a/master/src/main/java/org/apache/celeborn/service/deploy/master/slotsalloc/SlotsAllocator.java
b/master/src/main/java/org/apache/celeborn/service/deploy/master/slotsalloc/SlotsAllocator.java
index 4de3ed5c2..8f1cf5d4a 100644
---
a/master/src/main/java/org/apache/celeborn/service/deploy/master/slotsalloc/SlotsAllocator.java
+++
b/master/src/main/java/org/apache/celeborn/service/deploy/master/slotsalloc/SlotsAllocator.java
@@ -458,6 +458,7 @@ public class SlotsAllocator {
replicaIndex,
index ->
!(sameWorkerCandidates && index == selectedPrimaryIndex)
+ && canAssign(null, replicaWorkers.get(index),
availableStorageTypes)
&& satisfyRackAware(true, primaryWorker,
replicaWorkers.get(index)));
} else if (StorageInfo.localDiskAvailable(availableStorageTypes)) {
selectedReplicaIndex =
@@ -476,10 +477,7 @@ public class SlotsAllocator {
WorkerInfo replicaWorker = replicaWorkers.get(selectedReplicaIndex);
StorageInfo replicaStorageInfo =
- slotBudgets == null && shouldRackAware
- ? primaryStorageInfo
- : buildStorageInfo(
- replicaWorker, slotBudgets, workerDiskIndex,
availableStorageTypes);
+ buildStorageInfo(replicaWorker, slotBudgets, workerDiskIndex,
availableStorageTypes);
PartitionLocation replicaPartition =
createLocation(
partitionId,
diff --git
a/master/src/test/java/org/apache/celeborn/service/deploy/master/slotsalloc/SlotsAllocatorRackAwareSuiteJ.java
b/master/src/test/java/org/apache/celeborn/service/deploy/master/slotsalloc/SlotsAllocatorRackAwareSuiteJ.java
index ef13ef162..14b1d8df6 100644
---
a/master/src/test/java/org/apache/celeborn/service/deploy/master/slotsalloc/SlotsAllocatorRackAwareSuiteJ.java
+++
b/master/src/test/java/org/apache/celeborn/service/deploy/master/slotsalloc/SlotsAllocatorRackAwareSuiteJ.java
@@ -183,6 +183,95 @@ public class SlotsAllocatorRackAwareSuiteJ {
return workers;
}
+ @Test
+ public void offerSlotsRackAwareFallbackRequiresReplicaDisk() {
+ List<WorkerInfo> workers =
+ Arrays.asList(
+ prepareWorker("disk-worker", "/rack/r1", "/mnt/disk1"),
+ prepareWorker("diskless-worker", "/rack/r2", null));
+
+ Map<WorkerInfo, Tuple2<List<PartitionLocation>, List<PartitionLocation>>>
slots =
+ SlotsAllocator.offerSlots(
+ workers,
+ Collections.singletonList(0),
+ true,
+ true,
+ StorageInfo.LOCAL_DISK_MASK,
+ false,
+ 0,
+ SLOTS_ASSIGN_STRATEGY);
+
+ Assert.assertTrue(slots.isEmpty());
+ }
+
+ @Test
+ public void offerSlotsRackAwareFallbackSkipsDisklessReplica() {
+ WorkerInfo primaryWorker = prepareWorker("primary-worker", "/rack/r1",
"/mnt/disk1");
+ WorkerInfo disklessReplica = prepareWorker("diskless-replica", "/rack/r2",
null);
+ WorkerInfo diskReplica = prepareWorker("disk-replica", "/rack/r2",
"/mnt/disk2");
+ // Use separate candidate lists and two partitions so the diskless replica
is scanned regardless
+ // of the random initial replica index.
+ disklessReplica.nextInterruptionNotice_$eq(1L);
+ diskReplica.nextInterruptionNotice_$eq(2L);
+
+ Map<WorkerInfo, Tuple2<List<PartitionLocation>, List<PartitionLocation>>>
slots =
+ SlotsAllocator.offerSlots(
+ Arrays.asList(primaryWorker, disklessReplica, diskReplica),
+ Arrays.asList(0, 1),
+ true,
+ true,
+ StorageInfo.LOCAL_DISK_MASK,
+ true,
+ 0,
+ SLOTS_ASSIGN_STRATEGY);
+
+ Assert.assertFalse(slots.containsKey(disklessReplica));
+ Assert.assertEquals(2, slots.get(primaryWorker)._1.size());
+ Assert.assertEquals(2, slots.get(diskReplica)._2.size());
+ }
+
+ @Test
+ public void offerSlotsRackAwareFallbackUsesReplicaStorageInfo() {
+ List<WorkerInfo> workers =
+ Arrays.asList(
+ prepareWorker("worker1", "/rack/r1", "/mnt/disk1"),
+ prepareWorker("worker2", "/rack/r2", "/mnt/disk2"));
+
+ Map<WorkerInfo, Tuple2<List<PartitionLocation>, List<PartitionLocation>>>
slots =
+ SlotsAllocator.offerSlots(
+ workers,
+ Collections.singletonList(0),
+ true,
+ true,
+ StorageInfo.LOCAL_DISK_MASK,
+ false,
+ 0,
+ SLOTS_ASSIGN_STRATEGY);
+
+ Assert.assertEquals(2, slots.size());
+ workers.forEach(
+ worker -> {
+ List<PartitionLocation> locations = new ArrayList<>();
+ locations.addAll(slots.get(worker)._1);
+ locations.addAll(slots.get(worker)._2);
+ Assert.assertEquals(1, locations.size());
+ Assert.assertEquals(
+ worker.diskInfos().keySet().iterator().next(),
+ locations.get(0).getStorageInfo().getMountPoint());
+ });
+ }
+
+ private static WorkerInfo prepareWorker(String host, String rack, String
mountPoint) {
+ Map<String, DiskInfo> diskInfos = new HashMap<>();
+ if (mountPoint != null) {
+ // Keep availableSlots at zero to force allocation through the
best-effort fallback.
+ diskInfos.put(mountPoint, new DiskInfo(mountPoint, 1024L, 1, 1, 0));
+ }
+ WorkerInfo worker = new WorkerInfo(host, 1, 2, 3, 4, 5, diskInfos, null);
+ worker.networkLocation_$eq(rack);
+ return worker;
+ }
+
@Test
public void testRackAwareRoundRobinReplicaPatterns() {
for (Tuple2<Integer, List<WorkerInfo>> tuple :