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

cryptoe pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git


The following commit(s) were added to refs/heads/master by this push:
     new d1b99ecc204 minor: Prioritize load of unavailable segment (#20245)
d1b99ecc204 is described below

commit d1b99ecc20465f8ea84e810e290568822d99a09f
Author: Kashif Faraz <[email protected]>
AuthorDate: Wed Sep 9 20:06:17 2026 +0530

    minor: Prioritize load of unavailable segment (#20245)
    
    Description
    
    If a segment becomes unavailable after it has been added to the load queue 
of a historical, it might remain stuck in the queue while other segments are 
being loaded. Ideally, all unavailable segments should be prioritized.
    Changes
    
        Simplify book-keeping in ServerHolder by getting rid of method 
simplify() and treat REPLICATE and LOAD actions as distinct
        In StrategicSegmentAssigner.updateReplicasInTier(), check if an 
unavailable segment needs to be prioritized and change all in-flight REPLICATE 
actions on that segment to LOAD.
        Update LoadQueuePeon.getSegmentsInQueue to return a List instead of a 
Set and verify the result in Coordinator simulations
---
 .../druid/server/coordinator/ServerHolder.java     | 29 +++++---
 .../server/coordinator/duty/CloneHistoricals.java  |  2 +-
 .../coordinator/loading/HttpLoadQueuePeon.java     |  6 +-
 .../server/coordinator/loading/LoadQueuePeon.java  |  3 +-
 .../server/coordinator/loading/SegmentAction.java  |  1 -
 .../coordinator/loading/SegmentReplicaCount.java   | 16 +++++
 .../loading/StrategicSegmentAssigner.java          | 58 ++++++++++++++--
 .../server/coordinator/duty/RunRulesTest.java      |  6 +-
 .../coordinator/loading/TestLoadQueuePeon.java     |  6 +-
 .../simulate/CoordinatorSimulation.java            |  6 ++
 .../simulate/CoordinatorSimulationBaseTest.java    |  7 ++
 .../simulate/CoordinatorSimulationBuilder.java     |  9 +++
 .../coordinator/simulate/SegmentLoadingTest.java   | 80 ++++++++++++++++++++++
 13 files changed, 200 insertions(+), 29 deletions(-)

diff --git 
a/server/src/main/java/org/apache/druid/server/coordinator/ServerHolder.java 
b/server/src/main/java/org/apache/druid/server/coordinator/ServerHolder.java
index 0e137caed14..60593175d5f 100644
--- a/server/src/main/java/org/apache/druid/server/coordinator/ServerHolder.java
+++ b/server/src/main/java/org/apache/druid/server/coordinator/ServerHolder.java
@@ -167,7 +167,7 @@ public class ServerHolder implements 
Comparable<ServerHolder>
       }
 
       final SegmentAction action = holder.getAction();
-      addToQueuedSegments(holder.getSegment(), simplify(action));
+      addToQueuedSegments(holder.getSegment(), action);
       if (holder.getProfile() != null) {
         inFlightProfiles.put(holder.getSegment(), holder.getProfile());
       }
@@ -298,7 +298,6 @@ public class ServerHolder implements 
Comparable<ServerHolder>
    * <ul>
    * <li>Contains segments present in the queue when the current coordinator 
run started.</li>
    * <li>Contains segments added to the queue during the current run.</li>
-   * <li>Maps replicating segments to LOAD rather than REPLICATE for 
simplicity.</li>
    * <li>Does not contain segments whose actions were cancelled.</li>
    * </ul>
    */
@@ -350,7 +349,7 @@ public class ServerHolder implements 
Comparable<ServerHolder>
   {
     final List<DataSegment> loadingSegments = new ArrayList<>();
     queuedSegments.forEach((segment, action) -> {
-      if (action == SegmentAction.LOAD) {
+      if (action == SegmentAction.LOAD || action == SegmentAction.REPLICATE) {
         loadingSegments.add(segment);
       }
     });
@@ -376,7 +375,8 @@ public class ServerHolder implements 
Comparable<ServerHolder>
 
   public boolean isLoadingSegment(DataSegment segment)
   {
-    return getActionOnSegment(segment) == SegmentAction.LOAD;
+    final SegmentAction action = getActionOnSegment(segment);
+    return action == SegmentAction.LOAD || action == SegmentAction.REPLICATE;
   }
 
   public boolean isDroppingSegment(DataSegment segment)
@@ -431,7 +431,7 @@ public class ServerHolder implements 
Comparable<ServerHolder>
       ++totalAssignmentsInRun;
     }
 
-    addToQueuedSegments(segment, simplify(action));
+    addToQueuedSegments(segment, action);
     if (profile != null) {
       inFlightProfiles.put(segment, profile);
     }
@@ -442,7 +442,7 @@ public class ServerHolder implements 
Comparable<ServerHolder>
   {
     // Cancel only if the action is currently in queue
     final SegmentAction queuedAction = queuedSegments.get(segment);
-    if (queuedAction != simplify(action)) {
+    if (queuedAction != action) {
       return false;
     }
 
@@ -457,6 +457,18 @@ public class ServerHolder implements 
Comparable<ServerHolder>
     }
   }
 
+  /**
+   * Cancels a {@link SegmentAction#REPLICATE} or {@link SegmentAction#LOAD} if
+   * it is currently being performed on the given segment.
+   *
+   * @return true if the operation was cancelled successfully.
+   */
+  public boolean cancelLoad(DataSegment segment)
+  {
+    return cancelOperation(SegmentAction.REPLICATE, segment)
+           || cancelOperation(SegmentAction.LOAD, segment);
+  }
+
   /**
    * Returns the {@link PartialLoadProfile} for an in-flight load of {@code 
segment} on this server, or {@code null}
    * if there's no in-flight load or the in-flight load is a regular full-load 
(no profile). Used by the partial-load
@@ -497,11 +509,6 @@ public class ServerHolder implements 
Comparable<ServerHolder>
            || server.getType() == ServerType.INDEXER_EXECUTOR;
   }
 
-  private SegmentAction simplify(SegmentAction action)
-  {
-    return action == SegmentAction.REPLICATE ? SegmentAction.LOAD : action;
-  }
-
   private void addToQueuedSegments(DataSegment segment, SegmentAction action)
   {
     queuedSegments.put(segment, action);
diff --git 
a/server/src/main/java/org/apache/druid/server/coordinator/duty/CloneHistoricals.java
 
b/server/src/main/java/org/apache/druid/server/coordinator/duty/CloneHistoricals.java
index 902ed31a573..82441810634 100644
--- 
a/server/src/main/java/org/apache/druid/server/coordinator/duty/CloneHistoricals.java
+++ 
b/server/src/main/java/org/apache/druid/server/coordinator/duty/CloneHistoricals.java
@@ -185,7 +185,7 @@ public class CloneHistoricals implements CoordinatorDuty
   )
   {
     if (targetServer.isLoadingSegment(segment)) {
-      targetServer.cancelOperation(SegmentAction.LOAD, segment);
+      targetServer.cancelLoad(segment);
     } else if (loadQueueManager.dropSegment(segment, targetServer)) {
       params.getCoordinatorStats().add(
           Stats.Segments.DROPPED_FROM_CLONE,
diff --git 
a/server/src/main/java/org/apache/druid/server/coordinator/loading/HttpLoadQueuePeon.java
 
b/server/src/main/java/org/apache/druid/server/coordinator/loading/HttpLoadQueuePeon.java
index 438bea3a7bd..c51e2e703e7 100644
--- 
a/server/src/main/java/org/apache/druid/server/coordinator/loading/HttpLoadQueuePeon.java
+++ 
b/server/src/main/java/org/apache/druid/server/coordinator/loading/HttpLoadQueuePeon.java
@@ -704,11 +704,11 @@ public class HttpLoadQueuePeon implements LoadQueuePeon
   }
 
   @Override
-  public Set<SegmentHolder> getSegmentsInQueue()
+  public List<SegmentHolder> getSegmentsInQueue()
   {
-    final Set<SegmentHolder> segmentsInQueue;
+    final List<SegmentHolder> segmentsInQueue;
     synchronized (lock) {
-      segmentsInQueue = new HashSet<>(queuedSegments);
+      segmentsInQueue = new ArrayList<>(queuedSegments);
     }
     return segmentsInQueue;
   }
diff --git 
a/server/src/main/java/org/apache/druid/server/coordinator/loading/LoadQueuePeon.java
 
b/server/src/main/java/org/apache/druid/server/coordinator/loading/LoadQueuePeon.java
index 7590eae1e0d..bed24183cf2 100644
--- 
a/server/src/main/java/org/apache/druid/server/coordinator/loading/LoadQueuePeon.java
+++ 
b/server/src/main/java/org/apache/druid/server/coordinator/loading/LoadQueuePeon.java
@@ -23,6 +23,7 @@ import 
org.apache.druid.server.coordinator.stats.CoordinatorRunStats;
 import org.apache.druid.timeline.DataSegment;
 
 import javax.annotation.Nullable;
+import java.util.List;
 import java.util.Set;
 
 /**
@@ -36,7 +37,7 @@ public interface LoadQueuePeon
 
   Set<DataSegment> getSegmentsToLoad();
 
-  Set<SegmentHolder> getSegmentsInQueue();
+  List<SegmentHolder> getSegmentsInQueue();
 
   Set<DataSegment> getSegmentsToDrop();
 
diff --git 
a/server/src/main/java/org/apache/druid/server/coordinator/loading/SegmentAction.java
 
b/server/src/main/java/org/apache/druid/server/coordinator/loading/SegmentAction.java
index 7c5e983afcc..50c2e63488f 100644
--- 
a/server/src/main/java/org/apache/druid/server/coordinator/loading/SegmentAction.java
+++ 
b/server/src/main/java/org/apache/druid/server/coordinator/loading/SegmentAction.java
@@ -48,7 +48,6 @@ public enum SegmentAction
    *   <li>this action can be throttled by the {@code 
replicationThrottleLimit}</li>
    *   <li>it is given lower priority than LOAD on the load queue peon</li>
    * </ul>
-   * For all other purposes, REPLICATE is treated the same as LOAD.
    */
   REPLICATE(true),
 
diff --git 
a/server/src/main/java/org/apache/druid/server/coordinator/loading/SegmentReplicaCount.java
 
b/server/src/main/java/org/apache/druid/server/coordinator/loading/SegmentReplicaCount.java
index d2261080e5c..3c9bef3dcf4 100644
--- 
a/server/src/main/java/org/apache/druid/server/coordinator/loading/SegmentReplicaCount.java
+++ 
b/server/src/main/java/org/apache/druid/server/coordinator/loading/SegmentReplicaCount.java
@@ -37,6 +37,11 @@ public class SegmentReplicaCount
   private int movingTo;
   private int movingFrom;
 
+  /**
+   * Subset of {@link #loading}.
+   */
+  private int replicating;
+
   /**
    * Increments number of replicas loaded on historical servers.
    */
@@ -71,6 +76,7 @@ public class SegmentReplicaCount
   {
     switch (action) {
       case REPLICATE:
+        ++replicating;
       case LOAD:
         ++loading;
         break;
@@ -123,6 +129,15 @@ public class SegmentReplicaCount
     return loading;
   }
 
+  /**
+   * Number of replicas that are currently being loaded in the tier with action
+   * REPLICATE. Always less than or equal to {@link #loading()}.
+   */
+  int replicating()
+  {
+    return replicating;
+  }
+
   int moving()
   {
     return movingTo;
@@ -211,5 +226,6 @@ public class SegmentReplicaCount
     this.dropping += other.dropping;
     this.movingTo += other.movingTo;
     this.movingFrom += other.movingFrom;
+    this.replicating += other.replicating;
   }
 }
diff --git 
a/server/src/main/java/org/apache/druid/server/coordinator/loading/StrategicSegmentAssigner.java
 
b/server/src/main/java/org/apache/druid/server/coordinator/loading/StrategicSegmentAssigner.java
index ab95ed09cae..42cb7d3ab2c 100644
--- 
a/server/src/main/java/org/apache/druid/server/coordinator/loading/StrategicSegmentAssigner.java
+++ 
b/server/src/main/java/org/apache/druid/server/coordinator/loading/StrategicSegmentAssigner.java
@@ -43,6 +43,7 @@ import java.util.HashSet;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
+import java.util.Objects;
 import java.util.Set;
 import java.util.TreeSet;
 import java.util.stream.Collectors;
@@ -194,7 +195,7 @@ public class StrategicSegmentAssigner implements 
SegmentActionHandler
 
     if (serverA.isLoadingSegment(segment)) {
       // Cancel the load on serverA and load on serverB instead
-      if (serverA.cancelOperation(SegmentAction.LOAD, segment)) {
+      if (serverA.cancelLoad(segment)) {
         int loadedCountOnTier = replicaCountMap.get(segment.getId(), tier)
                                                .loadedNotDropping();
         if (loadedCountOnTier >= 1) {
@@ -592,9 +593,9 @@ public class StrategicSegmentAssigner implements 
SegmentActionHandler
       if (canceledOut.size() >= numToCancel) {
         break;
       }
-      // Try LOAD then REPLICATE; the queued action depends on whether this 
was a primary or a replica.
-      if (server.cancelOperation(SegmentAction.LOAD, segment)
-          || server.cancelOperation(SegmentAction.REPLICATE, segment)) {
+      // Try to cancel the load operation
+      // (either LOAD or REPLICATE, depending on whether this was a primary or 
a replica).
+      if (server.cancelLoad(segment)) {
         canceledOut.add(server);
       }
     }
@@ -658,8 +659,15 @@ public class StrategicSegmentAssigner implements 
SegmentActionHandler
     // clears the rule on the historical.
     final int replicasToRevert = requiredReplicas > 0 ? 
replicaCountOnTier.loadedWithPartialProfile() : 0;
 
+    final boolean shouldPrioritizeLoadOfUnavailableSegment =
+        replicaCountOnTier.loadedNotDropping() < 1
+        && replicaCountOnTier.replicating() >= 1;
+
     // Check if there is any action required on this tier
-    if (projectedReplicas == requiredReplicas && !shouldCancelMoves && 
replicasToRevert <= 0) {
+    if (projectedReplicas == requiredReplicas
+        && !shouldCancelMoves
+        && !shouldPrioritizeLoadOfUnavailableSegment
+        && replicasToRevert <= 0) {
       return 0;
     }
 
@@ -692,7 +700,9 @@ public class StrategicSegmentAssigner implements 
SegmentActionHandler
     if (projectedReplicas > requiredReplicas) {
       int replicaSurplus = projectedReplicas - requiredReplicas;
       int canceledLoads =
-          cancelOperations(SegmentAction.LOAD, replicaSurplus, segment, 
segmentStatus);
+          cancelOperations(SegmentAction.REPLICATE, replicaSurplus, segment, 
segmentStatus);
+      canceledLoads +=
+          cancelOperations(SegmentAction.LOAD, replicaSurplus - canceledLoads, 
segment, segmentStatus);
 
       int numReplicasToDrop = Math.min(replicaSurplus - canceledLoads, 
maxReplicasToDrop);
       if (numReplicasToDrop > 0) {
@@ -701,6 +711,13 @@ public class StrategicSegmentAssigner implements 
SegmentActionHandler
       }
     }
 
+    // If segment is unavailable, prioritize load by changing REPLICATE 
actions to LOAD
+    if (shouldPrioritizeLoadOfUnavailableSegment) {
+      for (ServerHolder server : 
segmentStatus.getServersPerforming(SegmentAction.REPLICATE)) {
+        prioritizeLoadOfUnavailableSegment(segment, server, null);
+      }
+    }
+
     // Release partial-load rules that no longer apply. Done last so the 
load/drop decisions above claim their
     // servers first: a replica that just picked up an action is no longer 
`isServingSegment`, so it is skipped here
     // and reverted on a later run if it is still around.
@@ -711,6 +728,8 @@ public class StrategicSegmentAssigner implements 
SegmentActionHandler
       }
     }
 
+
+
     return dropsQueuedOnTier;
   }
 
@@ -874,7 +893,7 @@ public class StrategicSegmentAssigner implements 
SegmentActionHandler
   private boolean dropBroadcastSegment(DataSegment segment, ServerHolder 
server)
   {
     if (server.isLoadingSegment(segment)) {
-      return server.cancelOperation(SegmentAction.LOAD, segment);
+      return server.cancelLoad(segment);
     } else if (server.isServingSegment(segment)) {
       return loadQueueManager.dropSegment(segment, server);
     } else {
@@ -1003,6 +1022,31 @@ public class StrategicSegmentAssigner implements 
SegmentActionHandler
     return numLoadsQueued;
   }
 
+  /**
+   * Tries to increase the load priority of the given unavailable segment (by
+   * changing the action from {@link SegmentAction#REPLICATE} to {@link 
SegmentAction#LOAD})
+   * if it is already present in the queue of the server.
+   *
+   * @return true only if the priority was increased successfully.
+   */
+  private boolean prioritizeLoadOfUnavailableSegment(
+      DataSegment segment,
+      ServerHolder server,
+      @Nullable PartialLoadProfile profile
+  )
+  {
+    return server.getActionOnSegment(segment) == SegmentAction.REPLICATE
+           && Objects.equals(fingerprintOf(profile), 
fingerprintOf(server.getProjectedProfile(segment)))
+           && server.cancelOperation(SegmentAction.REPLICATE, segment)
+           && loadQueueManager.loadSegment(segment, server, 
SegmentAction.LOAD, profile);
+  }
+
+  @Nullable
+  private static String fingerprintOf(@Nullable PartialLoadProfile profile)
+  {
+    return profile == null ? null : profile.fingerprint();
+  }
+
   private boolean loadSegment(DataSegment segment, ServerHolder server, 
@Nullable PartialLoadProfile profile)
   {
     final boolean assigned = loadQueueManager.loadSegment(segment, server, 
SegmentAction.LOAD, profile);
diff --git 
a/server/src/test/java/org/apache/druid/server/coordinator/duty/RunRulesTest.java
 
b/server/src/test/java/org/apache/druid/server/coordinator/duty/RunRulesTest.java
index 8cda9f8546a..67ffe445343 100644
--- 
a/server/src/test/java/org/apache/druid/server/coordinator/duty/RunRulesTest.java
+++ 
b/server/src/test/java/org/apache/druid/server/coordinator/duty/RunRulesTest.java
@@ -492,7 +492,7 @@ public class RunRulesTest
         )
     );
 
-    
EasyMock.expect(mockPeon.getSegmentsInQueue()).andReturn(Collections.emptySet()).anyTimes();
+    
EasyMock.expect(mockPeon.getSegmentsInQueue()).andReturn(Collections.emptyList()).anyTimes();
     
EasyMock.expect(mockPeon.getSegmentsMarkedToDrop()).andReturn(Collections.emptySet()).anyTimes();
     EasyMock.replay(mockPeon);
 
@@ -715,7 +715,7 @@ public class RunRulesTest
 
     LoadQueuePeon anotherMockPeon = EasyMock.createMock(LoadQueuePeon.class);
     
EasyMock.expect(anotherMockPeon.getSegmentsMarkedToDrop()).andReturn(Collections.emptySet()).anyTimes();
-    
EasyMock.expect(anotherMockPeon.getSegmentsInQueue()).andReturn(Collections.emptySet()).anyTimes();
+    
EasyMock.expect(anotherMockPeon.getSegmentsInQueue()).andReturn(Collections.emptyList()).anyTimes();
     
EasyMock.expect(anotherMockPeon.getSegmentsToLoad()).andReturn(Collections.emptySet()).anyTimes();
 
     EasyMock.replay(anotherMockPeon);
@@ -1227,7 +1227,7 @@ public class RunRulesTest
   {
     
EasyMock.expect(mockPeon.getSegmentsToLoad()).andReturn(Collections.emptySet()).anyTimes();
     
EasyMock.expect(mockPeon.getSegmentsMarkedToDrop()).andReturn(Collections.emptySet()).anyTimes();
-    
EasyMock.expect(mockPeon.getSegmentsInQueue()).andReturn(Collections.emptySet()).anyTimes();
+    
EasyMock.expect(mockPeon.getSegmentsInQueue()).andReturn(Collections.emptyList()).anyTimes();
     EasyMock.replay(mockPeon);
   }
 
diff --git 
a/server/src/test/java/org/apache/druid/server/coordinator/loading/TestLoadQueuePeon.java
 
b/server/src/test/java/org/apache/druid/server/coordinator/loading/TestLoadQueuePeon.java
index 87a09c6e6c0..1418e94dfef 100644
--- 
a/server/src/test/java/org/apache/druid/server/coordinator/loading/TestLoadQueuePeon.java
+++ 
b/server/src/test/java/org/apache/druid/server/coordinator/loading/TestLoadQueuePeon.java
@@ -23,7 +23,9 @@ import 
org.apache.druid.server.coordinator.stats.CoordinatorRunStats;
 import org.apache.druid.timeline.DataSegment;
 
 import javax.annotation.Nullable;
+import java.util.ArrayList;
 import java.util.Collections;
+import java.util.List;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ConcurrentSkipListSet;
@@ -138,9 +140,9 @@ public class TestLoadQueuePeon implements LoadQueuePeon
   }
 
   @Override
-  public Set<SegmentHolder> getSegmentsInQueue()
+  public List<SegmentHolder> getSegmentsInQueue()
   {
-    return queuedHolders;
+    return new ArrayList<>(queuedHolders);
   }
 
   @Override
diff --git 
a/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulation.java
 
b/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulation.java
index 5dd5e9b192d..8a6ff4b8205 100644
--- 
a/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulation.java
+++ 
b/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulation.java
@@ -22,6 +22,7 @@ package org.apache.druid.server.coordinator.simulate;
 import org.apache.druid.client.DruidServer;
 import org.apache.druid.java.util.metrics.MetricsVerifier;
 import org.apache.druid.server.coordinator.CoordinatorDynamicConfig;
+import org.apache.druid.server.coordinator.loading.SegmentHolder;
 import org.apache.druid.server.coordinator.rules.Rule;
 import org.apache.druid.timeline.DataSegment;
 
@@ -107,6 +108,11 @@ public interface CoordinatorSimulation
      */
     void loadQueuedSegments();
 
+    /**
+     * Gets the segments currently in the load queue of the given server.
+     */
+    List<SegmentHolder> getQueuedSegments(DruidServer server);
+
     /**
      * Finishes load of all the segments that were queued in the previous
      * coordinator run. Does not execute the respective callbacks on the 
coordinator.
diff --git 
a/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulationBaseTest.java
 
b/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulationBaseTest.java
index f42be68f8b9..9e8e16c8f25 100644
--- 
a/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulationBaseTest.java
+++ 
b/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulationBaseTest.java
@@ -26,6 +26,7 @@ import org.apache.druid.segment.TestDataSource;
 import org.apache.druid.server.coordination.ServerType;
 import org.apache.druid.server.coordinator.CoordinatorDynamicConfig;
 import org.apache.druid.server.coordinator.CreateDataSegments;
+import org.apache.druid.server.coordinator.loading.SegmentHolder;
 import 
org.apache.druid.server.coordinator.rules.ForeverBroadcastDistributionRule;
 import org.apache.druid.server.coordinator.rules.ForeverDropRule;
 import org.apache.druid.server.coordinator.rules.ForeverLoadRule;
@@ -126,6 +127,12 @@ public abstract class CoordinatorSimulationBaseTest 
implements
     sim.cluster().loadQueuedSegments();
   }
 
+  @Override
+  public List<SegmentHolder> getQueuedSegments(DruidServer server)
+  {
+    return sim.cluster().getQueuedSegments(server);
+  }
+
   @Override
   public void loadQueuedSegmentsSkipCallbacks()
   {
diff --git 
a/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulationBuilder.java
 
b/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulationBuilder.java
index 86ce9871aef..3484b7782d9 100644
--- 
a/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulationBuilder.java
+++ 
b/server/src/test/java/org/apache/druid/server/coordinator/simulate/CoordinatorSimulationBuilder.java
@@ -61,6 +61,7 @@ import 
org.apache.druid.server.coordinator.config.DruidCoordinatorConfig;
 import org.apache.druid.server.coordinator.config.HttpLoadQueuePeonConfig;
 import org.apache.druid.server.coordinator.duty.CoordinatorCustomDutyGroups;
 import org.apache.druid.server.coordinator.loading.LoadQueueTaskMaster;
+import org.apache.druid.server.coordinator.loading.SegmentHolder;
 import org.apache.druid.server.coordinator.loading.SegmentLoadQueueManager;
 import org.apache.druid.server.coordinator.rules.Rule;
 import org.apache.druid.server.http.BrokerDynamicConfigSyncer;
@@ -364,6 +365,14 @@ public class CoordinatorSimulationBuilder
       }
     }
 
+    @Override
+    public List<SegmentHolder> getQueuedSegments(DruidServer server)
+    {
+      return coordinator.getLoadManagementPeons()
+                        .get(server.getName())
+                        .getSegmentsInQueue();
+    }
+
     @Override
     public void removeServer(DruidServer server)
     {
diff --git 
a/server/src/test/java/org/apache/druid/server/coordinator/simulate/SegmentLoadingTest.java
 
b/server/src/test/java/org/apache/druid/server/coordinator/simulate/SegmentLoadingTest.java
index 430aec5b336..c831a3c4cb7 100644
--- 
a/server/src/test/java/org/apache/druid/server/coordinator/simulate/SegmentLoadingTest.java
+++ 
b/server/src/test/java/org/apache/druid/server/coordinator/simulate/SegmentLoadingTest.java
@@ -22,6 +22,8 @@ package org.apache.druid.server.coordinator.simulate;
 import org.apache.druid.client.DruidServer;
 import org.apache.druid.segment.TestDataSource;
 import org.apache.druid.server.coordinator.CoordinatorDynamicConfig;
+import org.apache.druid.server.coordinator.loading.SegmentAction;
+import org.apache.druid.server.coordinator.loading.SegmentHolder;
 import org.apache.druid.timeline.DataSegment;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeEach;
@@ -546,6 +548,84 @@ public class SegmentLoadingTest extends 
CoordinatorSimulationBaseTest
     Assertions.assertEquals(historicalT11.getCurrSize(), 
historicalT12.getCurrSize());
   }
 
+  @Test
+  public void testLoadQueuePrioritizesUnavailableSegment()
+  {
+    final CoordinatorSimulation sim =
+        CoordinatorSimulation.builder()
+                             .withServers(historicalT11, historicalT12)
+                             
.withDynamicConfig(withReplicationThrottleLimit(100))
+                             .withRules(datasource, Load.on(Tier.T1, 
2).forever())
+                             .build();
+
+    startSimulation(sim);
+
+    // All but last wiki segments are loaded on historicalT11
+    addSegments(segments);
+    for (int i = 0; i < 9; ++i) {
+      historicalT11.addDataSegment(segments.get(i));
+    }
+    final DataSegment unavailableSegment = segments.getLast();
+
+    runCoordinatorCycle();
+
+    // Verify that the load queue of historicalT12 has the unavailable segment 
first
+    final List<SegmentHolder> queuedSegments = 
getQueuedSegments(historicalT12);
+    Assertions.assertEquals(10, queuedSegments.size());
+
+    final SegmentHolder firstItemInQueue = queuedSegments.getFirst();
+    Assertions.assertEquals(unavailableSegment, firstItemInQueue.getSegment());
+    Assertions.assertEquals(SegmentAction.LOAD, firstItemInQueue.getAction());
+
+    for (int i = 1; i < queuedSegments.size(); ++i) {
+      Assertions.assertEquals(SegmentAction.REPLICATE, 
queuedSegments.get(i).getAction());
+    }
+  }
+
+  @Test
+  public void testSegmentMovesUpTheLoadQueueWhenItBecomesUnavailable()
+  {
+    final CoordinatorSimulation sim =
+        CoordinatorSimulation.builder()
+                             .withServers(historicalT11, historicalT12)
+                             
.withDynamicConfig(withReplicationThrottleLimit(100))
+                             .withRules(datasource, Load.on(Tier.T1, 
2).forever())
+                             .build();
+
+    startSimulation(sim);
+
+    // All wiki segments are loaded on historicalT11
+    addSegments(segments);
+    segments.forEach(historicalT11::addDataSegment);
+
+    runCoordinatorCycle();
+
+    // Verify that all segments in queue are for replication
+    final List<SegmentHolder> initialQueue = getQueuedSegments(historicalT12);
+    Assertions.assertEquals(10, initialQueue.size());
+    for (SegmentHolder holder : initialQueue) {
+      Assertions.assertEquals(SegmentAction.REPLICATE, holder.getAction());
+    }
+
+    // Now remove a couple of segments from historicalT11
+    final Set<DataSegment> missingSegments = Set.of(segments.get(1), 
segments.get(5));
+    missingSegments.forEach(segment -> 
historicalT11.removeDataSegment(segment.getId()));
+
+    runCoordinatorCycle();
+
+    // Verify that the missing segments have been moved up the queue
+    final List<SegmentHolder> updatedQueue = getQueuedSegments(historicalT12);
+    Assertions.assertEquals(10, updatedQueue.size());
+    for (int i = 0; i < 2; ++i) {
+      final SegmentHolder holder = updatedQueue.get(i);
+      Assertions.assertEquals(SegmentAction.LOAD, holder.getAction());
+      Assertions.assertTrue(missingSegments.contains(holder.getSegment()));
+    }
+    for (int i = 2; i < segments.size(); ++i) {
+      Assertions.assertEquals(SegmentAction.REPLICATE, 
updatedQueue.get(i).getAction());
+    }
+  }
+
   @Test
   public void testSegmentLoadingModes()
   {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to