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 :

Reply via email to