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]