This is an automated email from the ASF dual-hosted git repository.
capistrant 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 ed8efd5afdb fix: In coordinator duty loop - use an immutable snapshot
of retention rules per run cycle (#20063)
ed8efd5afdb is described below
commit ed8efd5afdb5c32aa14f4106b94636a2d6474427
Author: Lucas Capistrant <[email protected]>
AuthorDate: Mon Aug 24 08:58:03 2026 -0500
fix: In coordinator duty loop - use an immutable snapshot of retention
rules per run cycle (#20063)
---
docs/configuration/index.md | 2 +-
.../apache/druid/metadata/MetadataRuleManager.java | 14 +-
.../druid/metadata/MetadataRuleManagerConfig.java | 8 +-
.../druid/metadata/SQLMetadataRuleManager.java | 43 +--
.../druid/server/coordinator/DruidCoordinator.java | 7 +-
.../coordinator/DruidCoordinatorRuntimeParams.java | 27 ++
.../server/coordinator/duty/MetadataAction.java | 8 -
.../druid/server/coordinator/duty/RunRules.java | 21 +-
.../coordinator/duty/UnloadUnusedSegments.java | 25 +-
.../coordinator/rules/RetentionRulesSnapshot.java | 105 ++++++
.../druid/server/http/DataSourcesResource.java | 2 +-
.../apache/druid/server/http/RulesResource.java | 8 +-
.../druid/metadata/SQLMetadataRuleManagerTest.java | 87 ++++-
.../coordinator/BalanceSegmentsProfiler.java | 19 +-
.../server/coordinator/DruidCoordinatorTest.java | 30 +-
.../duty/RunRulesPartialLoadPlacementTest.java | 10 +-
.../server/coordinator/duty/RunRulesTest.java | 373 ++++++++-------------
.../coordinator/duty/UnloadUnusedSegmentsTest.java | 128 ++++---
.../server/coordinator/rules/LoadRuleTest.java | 42 +++
.../rules/RetentionRulesSnapshotTest.java | 168 ++++++++++
.../simulate/TestMetadataRuleManager.java | 30 +-
.../druid/server/http/DataSourcesResourceTest.java | 41 ++-
22 files changed, 772 insertions(+), 426 deletions(-)
diff --git a/docs/configuration/index.md b/docs/configuration/index.md
index ffe853cadf1..b4f676b4c0b 100644
--- a/docs/configuration/index.md
+++ b/docs/configuration/index.md
@@ -806,7 +806,7 @@ The following table shows the dynamic configuration
properties for the Coordinat
|`replicateAfterLoadTimeout`|Boolean flag for whether or not additional
replication is needed for segments that have failed to load due to the expiry
of `druid.coordinator.load.timeout`. If this is set to true, the Coordinator
will attempt to replicate the failed segment on a different historical server.
This helps improve the segment availability if there are a few slow Historicals
in the cluster. However, the slow Historical may still load the segment later
and the Coordinator may issu [...]
|`turboLoadingNodes`| Experimental. List of Historical servers to place in
turbo loading mode. These servers use a larger thread-pool to load segments
faster but at the cost of query performance. For servers specified in
`turboLoadingNodes`, `druid.coordinator.loadqueuepeon.http.batchSize` is
ignored and the coordinator uses the value of the respective
`numLoadingThreads` instead.<br/>Please use this config with caution. All
servers should eventually be removed from this list once the se [...]
|`cloneServers`| Experimental. Map from target Historical server to source
Historical server which should be cloned by the target. The target Historical
does not participate in regular segment assignment or balancing. Instead, the
Coordinator mirrors any segment assignment made to the source Historical onto
the target Historical, so that the target becomes an exact copy of the source.
Segments on the target Historical do not count towards replica counts either.
If the source disappears, [...]
-|`historicalTierAliases`|Map from a virtual tier name to the set of real
Historical tier names it expands to. When a load/drop rule references a virtual
alias tier, the Coordinator replaces it with its real tiers — each receiving
the full replica count independently. The alias key itself is never loaded to
directly. For example, `{"hot": ["hot_1", "hot_2"]}` causes a rule of `{"hot":
2}` to load 2 replicas on each of `hot_1` and `hot_2`; `hot` receives no direct
assignment. An alias valu [...]
+|`historicalTierAliases`|Map from a virtual tier name to the set of real
Historical tier names it expands to. When a load/drop rule references a virtual
alias tier, the Coordinator replaces it with its real tiers — each receiving
the full replica count independently. The alias key itself is never loaded to
directly. For example, `{"hot": ["hot_1", "hot_2"]}` causes a rule of `{"hot":
2}` to load 2 replicas on each of `hot_1` and `hot_2`; `hot` receives no direct
assignment. An alias valu [...]
##### Smart segment loading
diff --git
a/server/src/main/java/org/apache/druid/metadata/MetadataRuleManager.java
b/server/src/main/java/org/apache/druid/metadata/MetadataRuleManager.java
index ea2b6e7461f..036cdf7eef9 100644
--- a/server/src/main/java/org/apache/druid/metadata/MetadataRuleManager.java
+++ b/server/src/main/java/org/apache/druid/metadata/MetadataRuleManager.java
@@ -20,10 +20,10 @@
package org.apache.druid.metadata;
import org.apache.druid.audit.AuditInfo;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.rules.Rule;
import java.util.List;
-import java.util.Map;
/**
*/
@@ -35,11 +35,13 @@ public interface MetadataRuleManager
void poll();
- Map<String, List<Rule>> getAllRules();
-
- List<Rule> getRules(String dataSource);
-
- List<Rule> getRulesWithDefault(String dataSource);
+ /**
+ * Current snapshot of the rules of all datasources.
+ * <p>
+ * Using a single snapshot while performing an operation (such as a
Coordinator duty run) allows the steps within the
+ * operation to remain consistent with each other, even if the rules are
updated concurrently.
+ */
+ RetentionRulesSnapshot getRulesSnapshot();
boolean overrideRule(String dataSource, List<Rule> rulesConfig, AuditInfo
auditInfo);
diff --git
a/server/src/main/java/org/apache/druid/metadata/MetadataRuleManagerConfig.java
b/server/src/main/java/org/apache/druid/metadata/MetadataRuleManagerConfig.java
index ef06631dd17..8ade2d8b5e9 100644
---
a/server/src/main/java/org/apache/druid/metadata/MetadataRuleManagerConfig.java
+++
b/server/src/main/java/org/apache/druid/metadata/MetadataRuleManagerConfig.java
@@ -26,8 +26,14 @@ import org.joda.time.Period;
*/
public class MetadataRuleManagerConfig
{
+ /**
+ * Default value of {@code druid.manager.rules.defaultRule}, i.e. the
datasource against which
+ * the cluster-level default rules are stored when an operator has not
configured another name.
+ */
+ public static final String DEFAULT_RULE_NAME = "_default";
+
@JsonProperty
- private String defaultRule = "_default";
+ private String defaultRule = DEFAULT_RULE_NAME;
@JsonProperty
private Period pollDuration = new Period("PT1M");
diff --git
a/server/src/main/java/org/apache/druid/metadata/SQLMetadataRuleManager.java
b/server/src/main/java/org/apache/druid/metadata/SQLMetadataRuleManager.java
index dc1bde217d1..6fa3fd5891d 100644
--- a/server/src/main/java/org/apache/druid/metadata/SQLMetadataRuleManager.java
+++ b/server/src/main/java/org/apache/druid/metadata/SQLMetadataRuleManager.java
@@ -40,6 +40,7 @@ import
org.apache.druid.java.util.common.lifecycle.LifecycleStart;
import org.apache.druid.java.util.common.lifecycle.LifecycleStop;
import org.apache.druid.java.util.emitter.EmittingLogger;
import org.apache.druid.server.coordinator.rules.ForeverLoadRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.rules.Rule;
import org.joda.time.DateTime;
import org.skife.jdbi.v2.Handle;
@@ -48,7 +49,6 @@ import org.skife.jdbi.v2.Update;
import org.skife.jdbi.v2.tweak.HandleCallback;
import java.io.IOException;
-import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
@@ -129,7 +129,11 @@ public class SQLMetadataRuleManager implements
MetadataRuleManager
private final MetadataRuleManagerConfig config;
private final MetadataStorageTablesConfig dbTables;
private final IDBI dbi;
- private final AtomicReference<ImmutableMap<String, List<Rule>>> rules;
+ /**
+ * Snapshot of the rules read by the latest successful {@link #poll()}.
Rebuilt only on
+ * poll so that every reader shares the same instance.
+ */
+ private final AtomicReference<RetentionRulesSnapshot> rulesSnapshot;
private final AuditManager auditManager;
private final Object lock = new Object();
@@ -174,7 +178,7 @@ public class SQLMetadataRuleManager implements
MetadataRuleManager
"If specified, 'druid.manager.rules.defaultRule' must have a non-empty
value."
);
- this.rules = new AtomicReference<>(ImmutableMap.of());
+ this.rulesSnapshot = new AtomicReference<>(new
RetentionRulesSnapshot(Map.of(), defaultRule));
}
@Override
@@ -223,7 +227,7 @@ public class SQLMetadataRuleManager implements
MetadataRuleManager
if (currentStartOrder == -1) {
return;
}
- rules.set(ImmutableMap.of());
+ rulesSnapshot.set(new RetentionRulesSnapshot(Map.of(),
config.getDefaultRule()));
currentStartOrder = -1;
// This call cancels the periodic poll() task, scheduled in start().
exec.shutdownNow();
@@ -279,7 +283,7 @@ public class SQLMetadataRuleManager implements
MetadataRuleManager
final int newRuleCount =
newRules.values().stream().mapToInt(List::size).sum();
log.info("Polled and found [%d] rule(s) for [%d] datasource(s).",
newRuleCount, newRules.size());
- rules.set(newRules);
+ rulesSnapshot.set(new RetentionRulesSnapshot(newRules,
config.getDefaultRule()));
failStartTimeMs = 0;
}
catch (Exception e) {
@@ -298,30 +302,9 @@ public class SQLMetadataRuleManager implements
MetadataRuleManager
}
@Override
- public Map<String, List<Rule>> getAllRules()
- {
- return rules.get();
- }
-
- @Override
- public List<Rule> getRules(final String dataSource)
+ public RetentionRulesSnapshot getRulesSnapshot()
{
- List<Rule> retVal = rules.get().get(dataSource);
- return retVal == null ? new ArrayList<>() : retVal;
- }
-
- @Override
- public List<Rule> getRulesWithDefault(final String dataSource)
- {
- List<Rule> retVal = new ArrayList<>();
- Map<String, List<Rule>> theRules = rules.get();
- if (theRules.get(dataSource) != null) {
- retVal.addAll(theRules.get(dataSource));
- }
- if (theRules.get(config.getDefaultRule()) != null) {
- retVal.addAll(theRules.get(config.getDefaultRule()));
- }
- return retVal;
+ return rulesSnapshot.get();
}
@Override
@@ -336,7 +319,9 @@ public class SQLMetadataRuleManager implements
MetadataRuleManager
final String ruleString;
try {
ruleString = jsonMapper.writeValueAsString(newRules);
- if
(ruleString.equals(jsonMapper.writeValueAsString(rules.get().get(dataSource))))
{
+ // uses getAllRules over getOverrideRules to allow proper audit trail
when setting explicit [] override rules for
+ // a datasource with no existing overrides.
+ if
(ruleString.equals(jsonMapper.writeValueAsString(rulesSnapshot.get().getAllRules().get(dataSource))))
{
log.info("Retention rules unchanged for datasource[%s] with
rules[%s]", dataSource, ruleString);
return true;
}
diff --git
a/server/src/main/java/org/apache/druid/server/coordinator/DruidCoordinator.java
b/server/src/main/java/org/apache/druid/server/coordinator/DruidCoordinator.java
index dd12caf904e..e61df886763 100644
---
a/server/src/main/java/org/apache/druid/server/coordinator/DruidCoordinator.java
+++
b/server/src/main/java/org/apache/druid/server/coordinator/DruidCoordinator.java
@@ -543,8 +543,6 @@ public class DruidCoordinator
private List<CoordinatorDuty> makeHistoricalManagementDuties()
{
final MetadataAction.DeleteSegments deleteSegments =
this::markSegmentsAsUnused;
- final MetadataAction.GetDatasourceRules getRules
- = dataSource ->
metadataManager.rules().getRulesWithDefault(dataSource);
return ImmutableList.of(
new PrepareBalancerAndLoadQueues(
@@ -553,10 +551,10 @@ public class DruidCoordinator
balancerStrategyFactory,
serverInventoryView
),
- new RunRules(deleteSegments, getRules),
+ new RunRules(deleteSegments),
new UpdateReplicationStatus(),
new CollectSegmentStats(),
- new UnloadUnusedSegments(loadQueueManager, getRules),
+ new UnloadUnusedSegments(loadQueueManager),
new MarkOvershadowedSegmentsAsUnused(deleteSegments),
new MarkEternityTombstonesAsUnused(deleteSegments),
new BalanceSegments(config.getCoordinatorPeriod()),
@@ -736,6 +734,7 @@ public class DruidCoordinator
.builder()
.withDataSourcesSnapshot(dataSourcesSnapshot)
.withDynamicConfigs(metadataManager.configs().getCurrentDynamicConfig())
+
.withRetentionRulesSnapshot(metadataManager.rules().getRulesSnapshot())
.withCompactionConfig(metadataManager.configs().getCurrentCompactionConfig())
.build();
dutyGroup.run(params);
diff --git
a/server/src/main/java/org/apache/druid/server/coordinator/DruidCoordinatorRuntimeParams.java
b/server/src/main/java/org/apache/druid/server/coordinator/DruidCoordinatorRuntimeParams.java
index 576b2155ac7..a79bdc06579 100644
---
a/server/src/main/java/org/apache/druid/server/coordinator/DruidCoordinatorRuntimeParams.java
+++
b/server/src/main/java/org/apache/druid/server/coordinator/DruidCoordinatorRuntimeParams.java
@@ -27,6 +27,7 @@ import
org.apache.druid.server.coordinator.loading.SegmentLoadQueueManager;
import org.apache.druid.server.coordinator.loading.SegmentLoadingConfig;
import org.apache.druid.server.coordinator.loading.SegmentReplicationStatus;
import org.apache.druid.server.coordinator.loading.StrategicSegmentAssigner;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.stats.CoordinatorRunStats;
import org.apache.druid.server.coordinator.stats.Dimension;
import org.apache.druid.timeline.DataSegment;
@@ -51,6 +52,7 @@ public class DruidCoordinatorRuntimeParams
private final CoordinatorDynamicConfig coordinatorDynamicConfig;
private final DruidCompactionConfig compactionConfig;
private final SegmentLoadingConfig segmentLoadingConfig;
+ private final RetentionRulesSnapshot retentionRulesSnapshot;
private final CoordinatorRunStats stats;
private final BalancerStrategy balancerStrategy;
private final Set<String> broadcastDatasources;
@@ -63,6 +65,7 @@ public class DruidCoordinatorRuntimeParams
CoordinatorDynamicConfig coordinatorDynamicConfig,
DruidCompactionConfig compactionConfig,
SegmentLoadingConfig segmentLoadingConfig,
+ RetentionRulesSnapshot retentionRulesSnapshot,
CoordinatorRunStats stats,
BalancerStrategy balancerStrategy,
Set<String> broadcastDatasources
@@ -75,6 +78,7 @@ public class DruidCoordinatorRuntimeParams
this.coordinatorDynamicConfig = coordinatorDynamicConfig;
this.compactionConfig = compactionConfig;
this.segmentLoadingConfig = segmentLoadingConfig;
+ this.retentionRulesSnapshot = retentionRulesSnapshot;
this.stats = stats;
this.balancerStrategy = balancerStrategy;
this.broadcastDatasources = broadcastDatasources;
@@ -148,6 +152,15 @@ public class DruidCoordinatorRuntimeParams
return segmentLoadingConfig;
}
+ /**
+ * Retention rules of all datasources, snapshotted at the start of this run.
+ */
+ public RetentionRulesSnapshot getRetentionRulesSnapshot()
+ {
+ Preconditions.checkState(retentionRulesSnapshot != null,
"retentionRulesSnapshot must be set");
+ return retentionRulesSnapshot;
+ }
+
public CoordinatorRunStats getCoordinatorStats()
{
return stats;
@@ -184,6 +197,7 @@ public class DruidCoordinatorRuntimeParams
coordinatorDynamicConfig,
compactionConfig,
segmentLoadingConfig,
+ retentionRulesSnapshot,
stats,
balancerStrategy,
broadcastDatasources
@@ -200,6 +214,7 @@ public class DruidCoordinatorRuntimeParams
private CoordinatorDynamicConfig coordinatorDynamicConfig;
private DruidCompactionConfig compactionConfig;
private SegmentLoadingConfig segmentLoadingConfig;
+ private RetentionRulesSnapshot retentionRulesSnapshot;
private CoordinatorRunStats stats;
private BalancerStrategy balancerStrategy;
private Set<String> broadcastDatasources;
@@ -219,6 +234,7 @@ public class DruidCoordinatorRuntimeParams
CoordinatorDynamicConfig coordinatorDynamicConfig,
DruidCompactionConfig compactionConfig,
SegmentLoadingConfig segmentLoadingConfig,
+ RetentionRulesSnapshot retentionRulesSnapshot,
CoordinatorRunStats stats,
BalancerStrategy balancerStrategy,
Set<String> broadcastDatasources
@@ -231,6 +247,7 @@ public class DruidCoordinatorRuntimeParams
this.coordinatorDynamicConfig = coordinatorDynamicConfig;
this.compactionConfig = compactionConfig;
this.segmentLoadingConfig = segmentLoadingConfig;
+ this.retentionRulesSnapshot = retentionRulesSnapshot;
this.stats = stats;
this.balancerStrategy = balancerStrategy;
this.broadcastDatasources = broadcastDatasources;
@@ -252,6 +269,7 @@ public class DruidCoordinatorRuntimeParams
coordinatorDynamicConfig,
compactionConfig,
segmentLoadingConfig,
+ retentionRulesSnapshot,
stats,
balancerStrategy,
broadcastDatasources
@@ -346,6 +364,15 @@ public class DruidCoordinatorRuntimeParams
return this;
}
+ /**
+ * Sets the snapshot of retention rules to be used by all duties in this
run.
+ */
+ public Builder withRetentionRulesSnapshot(RetentionRulesSnapshot
retentionRulesSnapshot)
+ {
+ this.retentionRulesSnapshot = retentionRulesSnapshot;
+ return this;
+ }
+
public Builder withCompactionConfig(DruidCompactionConfig config)
{
this.compactionConfig = config;
diff --git
a/server/src/main/java/org/apache/druid/server/coordinator/duty/MetadataAction.java
b/server/src/main/java/org/apache/druid/server/coordinator/duty/MetadataAction.java
index dd6e886a4bb..a267b95aaac 100644
---
a/server/src/main/java/org/apache/druid/server/coordinator/duty/MetadataAction.java
+++
b/server/src/main/java/org/apache/druid/server/coordinator/duty/MetadataAction.java
@@ -19,10 +19,8 @@
package org.apache.druid.server.coordinator.duty;
-import org.apache.druid.server.coordinator.rules.Rule;
import org.apache.druid.timeline.SegmentId;
-import java.util.List;
import java.util.Set;
/**
@@ -36,10 +34,4 @@ public final class MetadataAction
{
int markSegmentsAsUnused(String datasource, Set<SegmentId> segmentIds);
}
-
- @FunctionalInterface
- public interface GetDatasourceRules
- {
- List<Rule> getRulesWithDefault(String dataSource);
- }
}
diff --git
a/server/src/main/java/org/apache/druid/server/coordinator/duty/RunRules.java
b/server/src/main/java/org/apache/druid/server/coordinator/duty/RunRules.java
index a03d7fe64bc..be008a24bce 100644
---
a/server/src/main/java/org/apache/druid/server/coordinator/duty/RunRules.java
+++
b/server/src/main/java/org/apache/druid/server/coordinator/duty/RunRules.java
@@ -28,6 +28,7 @@ import org.apache.druid.server.coordinator.DruidCluster;
import org.apache.druid.server.coordinator.DruidCoordinatorRuntimeParams;
import org.apache.druid.server.coordinator.loading.StrategicSegmentAssigner;
import org.apache.druid.server.coordinator.rules.BroadcastDistributionRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.rules.Rule;
import org.apache.druid.server.coordinator.stats.CoordinatorRunStats;
import org.apache.druid.server.coordinator.stats.Dimension;
@@ -56,15 +57,10 @@ public class RunRules implements CoordinatorDuty
private static final EmittingLogger log = new EmittingLogger(RunRules.class);
private final MetadataAction.DeleteSegments deleteHandler;
- private final MetadataAction.GetDatasourceRules ruleHandler;
- public RunRules(
- MetadataAction.DeleteSegments deleteHandler,
- MetadataAction.GetDatasourceRules ruleHandler
- )
+ public RunRules(MetadataAction.DeleteSegments deleteHandler)
{
this.deleteHandler = deleteHandler;
- this.ruleHandler = ruleHandler;
}
@Override
@@ -83,6 +79,7 @@ public class RunRules implements CoordinatorDuty
final Set<DataSegment> overshadowed =
params.getDataSourcesSnapshot().getOvershadowedSegments();
final StrategicSegmentAssigner segmentAssigner =
params.getSegmentAssigner();
+ final RetentionRulesSnapshot rulesSnapshot =
params.getRetentionRulesSnapshot();
final DateTime now = DateTimes.nowUtc();
final Object2IntOpenHashMap<String> datasourceToSegmentsWithNoRule = new
Object2IntOpenHashMap<>();
@@ -95,7 +92,7 @@ public class RunRules implements CoordinatorDuty
}
// Find and apply matching rule
- List<Rule> rules =
ruleHandler.getRulesWithDefault(segment.getDataSource());
+ List<Rule> rules =
rulesSnapshot.getEffectiveRules(segment.getDataSource());
boolean foundMatchingRule = false;
for (Rule rule : rules) {
if (rule.appliesTo(segment, now)) {
@@ -163,7 +160,7 @@ public class RunRules implements CoordinatorDuty
{
return
params.getDataSourcesSnapshot().getDataSourcesMap().values().stream()
.map(ImmutableDruidDataSource::getName)
- .filter(this::isBroadcastDatasource)
+ .filter(datasource -> isBroadcastDatasource(datasource,
params.getRetentionRulesSnapshot()))
.collect(Collectors.toSet());
}
@@ -175,10 +172,10 @@ public class RunRules implements CoordinatorDuty
* <li>Are unloaded if unused, even from realtime servers</li>
* </ul>
*/
- private boolean isBroadcastDatasource(String datasource)
+ private boolean isBroadcastDatasource(String datasource,
RetentionRulesSnapshot rulesSnapshot)
{
- return ruleHandler.getRulesWithDefault(datasource)
- .stream()
- .anyMatch(rule -> rule instanceof
BroadcastDistributionRule);
+ return rulesSnapshot.getEffectiveRules(datasource)
+ .stream()
+ .anyMatch(rule -> rule instanceof
BroadcastDistributionRule);
}
}
diff --git
a/server/src/main/java/org/apache/druid/server/coordinator/duty/UnloadUnusedSegments.java
b/server/src/main/java/org/apache/druid/server/coordinator/duty/UnloadUnusedSegments.java
index 761e3383ede..243dc726fbb 100644
---
a/server/src/main/java/org/apache/druid/server/coordinator/duty/UnloadUnusedSegments.java
+++
b/server/src/main/java/org/apache/druid/server/coordinator/duty/UnloadUnusedSegments.java
@@ -26,6 +26,7 @@ import
org.apache.druid.server.coordinator.DruidCoordinatorRuntimeParams;
import org.apache.druid.server.coordinator.ServerHolder;
import org.apache.druid.server.coordinator.loading.SegmentLoadQueueManager;
import org.apache.druid.server.coordinator.rules.BroadcastDistributionRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.stats.CoordinatorRunStats;
import org.apache.druid.server.coordinator.stats.Stats;
import org.apache.druid.timeline.DataSegment;
@@ -43,14 +44,9 @@ public class UnloadUnusedSegments implements CoordinatorDuty
private static final Logger log = new Logger(UnloadUnusedSegments.class);
private final SegmentLoadQueueManager loadQueueManager;
- private final MetadataAction.GetDatasourceRules ruleHandler;
- public UnloadUnusedSegments(
- SegmentLoadQueueManager loadQueueManager,
- MetadataAction.GetDatasourceRules ruleHandler
- )
+ public UnloadUnusedSegments(SegmentLoadQueueManager loadQueueManager)
{
- this.ruleHandler = ruleHandler;
this.loadQueueManager = loadQueueManager;
}
@@ -89,7 +85,7 @@ public class UnloadUnusedSegments implements CoordinatorDuty
final AtomicInteger numQueuedDrops = new AtomicInteger(0);
final ImmutableDruidServer server = serverHolder.getServer();
for (ImmutableDruidDataSource dataSource : server.getDataSources()) {
- if (shouldSkipUnload(serverHolder, dataSource.getName(),
broadcastStatusByDatasource)) {
+ if (shouldSkipUnload(serverHolder, dataSource.getName(),
broadcastStatusByDatasource, params.getRetentionRulesSnapshot())) {
continue;
}
@@ -122,7 +118,7 @@ public class UnloadUnusedSegments implements CoordinatorDuty
{
final AtomicInteger cancelledOperations = new AtomicInteger(0);
server.getQueuedSegments().forEach((segment, action) -> {
- if (shouldSkipUnload(server, segment.getDataSource(),
broadcastStatusByDatasource)) {
+ if (shouldSkipUnload(server, segment.getDataSource(),
broadcastStatusByDatasource, params.getRetentionRulesSnapshot())) {
// do nothing
} else if (params.isUsedSegment(segment)) {
// do nothing
@@ -147,21 +143,22 @@ public class UnloadUnusedSegments implements
CoordinatorDuty
private boolean shouldSkipUnload(
ServerHolder server,
String dataSource,
- Map<String, Boolean> broadcastStatusByDatasource
+ Map<String, Boolean> broadcastStatusByDatasource,
+ RetentionRulesSnapshot rulesSnapshot
)
{
boolean isBroadcastDatasource = broadcastStatusByDatasource
- .computeIfAbsent(dataSource, this::isBroadcastDatasource);
+ .computeIfAbsent(dataSource, ds -> isBroadcastDatasource(ds,
rulesSnapshot));
return server.isRealtimeServer() && !isBroadcastDatasource;
}
/**
* A datasource is considered a broadcast datasource if it has even one
broadcast rule.
*/
- private boolean isBroadcastDatasource(String datasource)
+ private boolean isBroadcastDatasource(String datasource,
RetentionRulesSnapshot rulesSnapshot)
{
- return ruleHandler.getRulesWithDefault(datasource)
- .stream()
- .anyMatch(rule -> rule instanceof
BroadcastDistributionRule);
+ return rulesSnapshot.getEffectiveRules(datasource)
+ .stream()
+ .anyMatch(rule -> rule instanceof
BroadcastDistributionRule);
}
}
diff --git
a/server/src/main/java/org/apache/druid/server/coordinator/rules/RetentionRulesSnapshot.java
b/server/src/main/java/org/apache/druid/server/coordinator/rules/RetentionRulesSnapshot.java
new file mode 100644
index 00000000000..5e6ec32d0ab
--- /dev/null
+++
b/server/src/main/java/org/apache/druid/server/coordinator/rules/RetentionRulesSnapshot.java
@@ -0,0 +1,105 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.druid.server.coordinator.rules;
+
+import com.google.common.collect.Maps;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Immutable snapshot of the retention rules of every datasource.
+ */
+public class RetentionRulesSnapshot
+{
+ /**
+ * Override rules of each datasource. Cluster level rules are present in
this map against the
+ * default datasource, which is where {@link #clusterDefaultRules} comes
from.
+ */
+ private final Map<String, List<Rule>> datasourceToRules;
+ /**
+ * Override rules of each datasource, already concatenated with {@link
#clusterDefaultRules}.
+ * Contains an entry only for datasources that have override rules and are
not the default
+ * datasource itself, so that everything else can share the {@link
#clusterDefaultRules} instance.
+ */
+ private final Map<String, List<Rule>> datasourceToEffectiveRules;
+ private final List<Rule> clusterDefaultRules;
+
+ /**
+ * @param datasourceToRules Rules configured for each datasource,
including the entry for
+ * {@code defaultDatasourceName}.
+ * @param defaultDatasourceName Name of the datasource whose rules serve as
the cluster
+ * defaults, i.e. {@code
druid.manager.rules.defaultRule}. The
+ * cluster defaults are empty if it has no
entry in
+ * {@code datasourceToRules}.
+ */
+ public RetentionRulesSnapshot(Map<String, List<Rule>> datasourceToRules,
String defaultDatasourceName)
+ {
+ this.clusterDefaultRules =
List.copyOf(datasourceToRules.getOrDefault(defaultDatasourceName, List.of()));
+
+ // Copy the rule lists as well as the map spine, so that a caller still
holding one of
+ // the source lists cannot mutate this snapshot. The effective rules of
each datasource are
+ // resolved here rather than in getEffectiveRules(), which is called once
per used segment.
+ final Map<String, List<Rule>> rules =
Maps.newHashMapWithExpectedSize(datasourceToRules.size());
+ final Map<String, List<Rule>> effectiveRules =
Maps.newHashMapWithExpectedSize(datasourceToRules.size());
+ datasourceToRules.forEach((datasource, overrideRules) -> {
+ rules.put(datasource, List.copyOf(overrideRules));
+ // The default datasource is skipped so that its own rules are not
appended to themselves.
+ if (!overrideRules.isEmpty() &&
!datasource.equals(defaultDatasourceName)) {
+ final List<Rule> combinedRules = new ArrayList<>(overrideRules.size()
+ this.clusterDefaultRules.size());
+ combinedRules.addAll(overrideRules);
+ combinedRules.addAll(this.clusterDefaultRules);
+ effectiveRules.put(datasource,
Collections.unmodifiableList(combinedRules));
+ }
+ });
+ this.datasourceToRules = Map.copyOf(rules);
+ this.datasourceToEffectiveRules = Map.copyOf(effectiveRules);
+ }
+
+ /**
+ * Return all rules that exist in the cluster.
+ */
+ public Map<String, List<Rule>> getAllRules()
+ {
+ return datasourceToRules;
+ }
+
+ /**
+ * Override rules configured for this datasource, excluding the cluster
defaults.
+ * <p>
+ * No cluster defaults are appended, so a datasource with no overrides
returns an empty
+ * list. Use {@link #getEffectiveRules} to get the rules that actually apply
to its segments.
+ */
+ public List<Rule> getOverrideRules(String datasource)
+ {
+ return datasourceToRules.getOrDefault(datasource, List.of());
+ }
+
+ /**
+ * All retention rules applicable to segments of this datasource.
+ * The returned list contains the override rules specified for the
datasource followed by the cluster default rules.
+ */
+ public List<Rule> getEffectiveRules(String datasource)
+ {
+ return datasourceToEffectiveRules.getOrDefault(datasource,
clusterDefaultRules);
+ }
+}
diff --git
a/server/src/main/java/org/apache/druid/server/http/DataSourcesResource.java
b/server/src/main/java/org/apache/druid/server/http/DataSourcesResource.java
index 31281842923..859ed3498ae 100644
--- a/server/src/main/java/org/apache/druid/server/http/DataSourcesResource.java
+++ b/server/src/main/java/org/apache/druid/server/http/DataSourcesResource.java
@@ -923,7 +923,7 @@ public class DataSourcesResource
)
{
try {
- final List<Rule> rules =
metadataRuleManager.getRulesWithDefault(dataSourceName);
+ final List<Rule> rules =
metadataRuleManager.getRulesSnapshot().getEffectiveRules(dataSourceName);
final Interval theInterval = Intervals.of(interval);
final SegmentDescriptor descriptor = new SegmentDescriptor(theInterval,
version, partitionNumber);
final DateTime now = DateTimes.nowUtc();
diff --git
a/server/src/main/java/org/apache/druid/server/http/RulesResource.java
b/server/src/main/java/org/apache/druid/server/http/RulesResource.java
index aec12008d3a..2f9f64b89b8 100644
--- a/server/src/main/java/org/apache/druid/server/http/RulesResource.java
+++ b/server/src/main/java/org/apache/druid/server/http/RulesResource.java
@@ -27,6 +27,7 @@ import org.apache.druid.audit.AuditInfo;
import org.apache.druid.audit.AuditManager;
import org.apache.druid.java.util.common.Intervals;
import org.apache.druid.metadata.MetadataRuleManager;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.rules.Rule;
import org.apache.druid.server.http.security.RulesResourceFilter;
import org.apache.druid.server.http.security.StateResourceFilter;
@@ -71,7 +72,7 @@ public class RulesResource
@ResourceFilters(StateResourceFilter.class)
public Response getRules()
{
- return Response.ok(databaseRuleManager.getAllRules()).build();
+ return
Response.ok(databaseRuleManager.getRulesSnapshot().getAllRules()).build();
}
@GET
@@ -83,11 +84,12 @@ public class RulesResource
@QueryParam("full") final String full
)
{
+ final RetentionRulesSnapshot rulesSnapshot =
databaseRuleManager.getRulesSnapshot();
if (full != null) {
- return
Response.ok(databaseRuleManager.getRulesWithDefault(dataSourceName))
+ return Response.ok(rulesSnapshot.getEffectiveRules(dataSourceName))
.build();
}
- return Response.ok(databaseRuleManager.getRules(dataSourceName))
+ return Response.ok(rulesSnapshot.getOverrideRules(dataSourceName))
.build();
}
diff --git
a/server/src/test/java/org/apache/druid/metadata/SQLMetadataRuleManagerTest.java
b/server/src/test/java/org/apache/druid/metadata/SQLMetadataRuleManagerTest.java
index c4ef6c56240..9d390a800f3 100644
---
a/server/src/test/java/org/apache/druid/metadata/SQLMetadataRuleManagerTest.java
+++
b/server/src/test/java/org/apache/druid/metadata/SQLMetadataRuleManagerTest.java
@@ -38,6 +38,7 @@ import org.apache.druid.server.audit.AuditSerdeHelper;
import org.apache.druid.server.audit.SQLAuditManager;
import org.apache.druid.server.audit.SQLAuditManagerConfig;
import org.apache.druid.server.coordinator.rules.IntervalLoadRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.rules.Rule;
import org.apache.druid.server.metrics.NoopServiceEmitter;
import org.apache.druid.timeline.DataSegment;
@@ -110,12 +111,70 @@ public class SQLMetadataRuleManagerTest
);
ruleManager.overrideRule(TestDataSource.WIKI, rules,
createAuditInfo("override rule"));
// New rule should be be reflected in the in memory rules map immediately
after being set by user
- Map<String, List<Rule>> allRules = ruleManager.getAllRules();
+ Map<String, List<Rule>> allRules =
ruleManager.getRulesSnapshot().getAllRules();
Assertions.assertEquals(1, allRules.size());
Assertions.assertEquals(1, allRules.get(TestDataSource.WIKI).size());
Assertions.assertEquals(rules.get(0),
allRules.get(TestDataSource.WIKI).get(0));
}
+ @Test
+ public void testGetRulesSnapshot()
+ {
+ // Creates the cluster default rule
+ ruleManager.start();
+
+ final List<Rule> wikiRules = Collections.singletonList(
+ new IntervalLoadRule(
+ Intervals.of("2015-01-01/2015-02-01"),
+ ImmutableMap.of(DruidServer.DEFAULT_TIER,
DruidServer.DEFAULT_NUM_REPLICANTS),
+ null
+ )
+ );
+ ruleManager.overrideRule(TestDataSource.WIKI, wikiRules,
createAuditInfo("override rule"));
+
+ final List<Rule> defaultRules =
ruleManager.getRulesSnapshot().getAllRules().get(managerConfig.getDefaultRule());
+ final RetentionRulesSnapshot snapshot = ruleManager.getRulesSnapshot();
+
+ // A datasource with rules of its own gets them ahead of the cluster
defaults
+ Assertions.assertEquals(
+ ImmutableList.builder().addAll(wikiRules).addAll(defaultRules).build(),
+ snapshot.getEffectiveRules(TestDataSource.WIKI)
+ );
+
+ // A datasource with no rules of its own gets only the cluster defaults
+ Assertions.assertEquals(defaultRules,
snapshot.getEffectiveRules(TestDataSource.KOALA));
+ }
+
+ @Test
+ public void testGetRulesSnapshotIsUnaffectedByLaterOverride()
+ {
+ ruleManager.start();
+
+ final List<Rule> originalRules = Collections.singletonList(
+ new IntervalLoadRule(
+ Intervals.of("2015-01-01/2015-02-01"),
+ ImmutableMap.of(DruidServer.DEFAULT_TIER, 1),
+ null
+ )
+ );
+ ruleManager.overrideRule(TestDataSource.WIKI, originalRules,
createAuditInfo("override rule"));
+
+ final RetentionRulesSnapshot snapshot = ruleManager.getRulesSnapshot();
+
+ final List<Rule> updatedRules = Collections.singletonList(
+ new IntervalLoadRule(
+ Intervals.of("2015-01-01/2015-02-01"),
+ ImmutableMap.of(DruidServer.DEFAULT_TIER, 2),
+ null
+ )
+ );
+ ruleManager.overrideRule(TestDataSource.WIKI, updatedRules,
createAuditInfo("override rule"));
+
+ // The manager has picked up the new rules, but the snapshot still serves
the old ones
+ Assertions.assertEquals(updatedRules.get(0),
ruleManager.getRulesSnapshot().getEffectiveRules(TestDataSource.WIKI).get(0));
+ Assertions.assertEquals(originalRules.get(0),
snapshot.getEffectiveRules(TestDataSource.WIKI).get(0));
+ }
+
@Test
public void testAuditEntryNotCreatedWhenRetentionRulesUnchanged()
{
@@ -223,7 +282,7 @@ public class SQLMetadataRuleManagerTest
// fetch rules from metadata storage
ruleManager.poll();
- Assertions.assertEquals(rules, ruleManager.getRules(TestDataSource.WIKI));
+ Assertions.assertEquals(rules,
ruleManager.getRulesSnapshot().getOverrideRules(TestDataSource.WIKI));
// verify audit entry is created
List<AuditEntry> auditEntries =
auditManager.fetchAuditHistory(TestDataSource.WIKI, "rules", null);
@@ -256,8 +315,8 @@ public class SQLMetadataRuleManagerTest
// fetch rules from metadata storage
ruleManager.poll();
- Assertions.assertEquals(rules, ruleManager.getRules(TestDataSource.WIKI));
- Assertions.assertEquals(rules, ruleManager.getRules("test_dataSource2"));
+ Assertions.assertEquals(rules,
ruleManager.getRulesSnapshot().getOverrideRules(TestDataSource.WIKI));
+ Assertions.assertEquals(rules,
ruleManager.getRulesSnapshot().getOverrideRules("test_dataSource2"));
// test fetch audit entries
List<AuditEntry> auditEntries = auditManager.fetchAuditHistory("rules",
null);
@@ -285,7 +344,7 @@ public class SQLMetadataRuleManagerTest
// Verify that the rule was added
ruleManager.poll();
- Map<String, List<Rule>> allRules = ruleManager.getAllRules();
+ Map<String, List<Rule>> allRules =
ruleManager.getRulesSnapshot().getAllRules();
Assertions.assertEquals(1, allRules.size());
Assertions.assertEquals(1, allRules.get(TestDataSource.WIKI).size());
@@ -294,7 +353,7 @@ public class SQLMetadataRuleManagerTest
// Verify that rule was deleted
ruleManager.poll();
- allRules = ruleManager.getAllRules();
+ allRules = ruleManager.getRulesSnapshot().getAllRules();
Assertions.assertEquals(0, allRules.size());
}
@@ -312,7 +371,7 @@ public class SQLMetadataRuleManagerTest
// Verify that rule was added
ruleManager.poll();
- Map<String, List<Rule>> allRules = ruleManager.getAllRules();
+ Map<String, List<Rule>> allRules =
ruleManager.getRulesSnapshot().getAllRules();
Assertions.assertEquals(1, allRules.size());
Assertions.assertEquals(1, allRules.get(TestDataSource.WIKI).size());
@@ -322,7 +381,7 @@ public class SQLMetadataRuleManagerTest
// Verify that rule was not deleted
ruleManager.poll();
- allRules = ruleManager.getAllRules();
+ allRules = ruleManager.getRulesSnapshot().getAllRules();
Assertions.assertEquals(1, allRules.size());
Assertions.assertEquals(1, allRules.get(TestDataSource.WIKI).size());
}
@@ -341,7 +400,7 @@ public class SQLMetadataRuleManagerTest
// Verify that rule was added
ruleManager.poll();
- Map<String, List<Rule>> allRules = ruleManager.getAllRules();
+ Map<String, List<Rule>> allRules =
ruleManager.getRulesSnapshot().getAllRules();
Assertions.assertEquals(1, allRules.size());
Assertions.assertEquals(1, allRules.get(TestDataSource.WIKI).size());
@@ -368,7 +427,7 @@ public class SQLMetadataRuleManagerTest
// Verify that rule was not deleted
ruleManager.poll();
- allRules = ruleManager.getAllRules();
+ allRules = ruleManager.getRulesSnapshot().getAllRules();
Assertions.assertEquals(1, allRules.size());
Assertions.assertEquals(1, allRules.get(TestDataSource.WIKI).size());
}
@@ -380,16 +439,16 @@ public class SQLMetadataRuleManagerTest
ruleManager.start();
// Verify the default rule
ruleManager.poll();
- Map<String, List<Rule>> allRules = ruleManager.getAllRules();
+ Map<String, List<Rule>> allRules =
ruleManager.getRulesSnapshot().getAllRules();
Assertions.assertEquals(1, allRules.size());
- Assertions.assertEquals(1, allRules.get("_default").size());
+ Assertions.assertEquals(1,
allRules.get(managerConfig.getDefaultRule()).size());
// Delete everything
ruleManager.removeRulesForEmptyDatasourcesOlderThan(System.currentTimeMillis());
// Verify the default rule was not deleted
ruleManager.poll();
- allRules = ruleManager.getAllRules();
+ allRules = ruleManager.getRulesSnapshot().getAllRules();
Assertions.assertEquals(1, allRules.size());
- Assertions.assertEquals(1, allRules.get("_default").size());
+ Assertions.assertEquals(1,
allRules.get(managerConfig.getDefaultRule()).size());
}
@AfterEach
diff --git
a/server/src/test/java/org/apache/druid/server/coordinator/BalanceSegmentsProfiler.java
b/server/src/test/java/org/apache/druid/server/coordinator/BalanceSegmentsProfiler.java
index 277ef792b09..08ac01cd37e 100644
---
a/server/src/test/java/org/apache/druid/server/coordinator/BalanceSegmentsProfiler.java
+++
b/server/src/test/java/org/apache/druid/server/coordinator/BalanceSegmentsProfiler.java
@@ -27,13 +27,14 @@ import org.apache.druid.client.ImmutableDruidServerTests;
import org.apache.druid.java.util.common.DateTimes;
import org.apache.druid.java.util.emitter.EmittingLogger;
import org.apache.druid.java.util.emitter.service.ServiceEmitter;
-import org.apache.druid.metadata.MetadataRuleManager;
+import org.apache.druid.metadata.MetadataRuleManagerConfig;
import org.apache.druid.server.coordinator.duty.BalanceSegments;
import org.apache.druid.server.coordinator.duty.RunRules;
import org.apache.druid.server.coordinator.loading.LoadQueuePeon;
import org.apache.druid.server.coordinator.loading.SegmentLoadQueueManager;
import org.apache.druid.server.coordinator.loading.TestLoadQueuePeon;
import org.apache.druid.server.coordinator.rules.PeriodLoadRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.rules.Rule;
import org.apache.druid.timeline.DataSegment;
import org.apache.druid.timeline.partition.NoneShardSpec;
@@ -47,6 +48,7 @@ import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
+import java.util.Map;
/**
* TODO convert benchmarks to JMH
@@ -59,7 +61,6 @@ public class BalanceSegmentsProfiler
private ImmutableDruidServer druidServer2;
List<DataSegment> segments = new ArrayList<>();
ServiceEmitter emitter;
- MetadataRuleManager manager;
PeriodLoadRule loadRule = new PeriodLoadRule(new Period("P5000Y"), null,
ImmutableMap.of("normal", 3), null);
List<Rule> rules = ImmutableList.of(loadRule);
@@ -71,7 +72,6 @@ public class BalanceSegmentsProfiler
druidServer2 = EasyMock.createMock(ImmutableDruidServer.class);
emitter = EasyMock.createMock(ServiceEmitter.class);
EmittingLogger.registerEmitter(emitter);
- manager = EasyMock.createMock(MetadataRuleManager.class);
}
public void bigProfiler()
@@ -79,11 +79,6 @@ public class BalanceSegmentsProfiler
Stopwatch watch = Stopwatch.createUnstarted();
int numSegments = 55000;
int numServers = 50;
- EasyMock.expect(manager.getAllRules()).andReturn(ImmutableMap.of("test",
rules)).anyTimes();
-
EasyMock.expect(manager.getRules(EasyMock.anyObject())).andReturn(rules).anyTimes();
-
EasyMock.expect(manager.getRulesWithDefault(EasyMock.anyObject())).andReturn(rules).anyTimes();
- EasyMock.replay(manager);
-
List<ServerHolder> serverHolderList = new ArrayList<>();
List<DataSegment> segments = new ArrayList<>();
for (int i = 0; i < numSegments; i++) {
@@ -130,6 +125,12 @@ public class BalanceSegmentsProfiler
.builder()
.withDruidCluster(druidCluster)
.withUsedSegments(segments)
+ .withRetentionRulesSnapshot(
+ new RetentionRulesSnapshot(
+ Map.of(MetadataRuleManagerConfig.DEFAULT_RULE_NAME, rules),
+ MetadataRuleManagerConfig.DEFAULT_RULE_NAME
+ )
+ )
.withDynamicConfigs(
CoordinatorDynamicConfig
.builder()
@@ -142,7 +143,7 @@ public class BalanceSegmentsProfiler
.build();
BalanceSegments tester = new BalanceSegments(Duration.standardMinutes(1));
- RunRules runner = new RunRules((ds, set) -> set.size(),
manager::getRulesWithDefault);
+ RunRules runner = new RunRules((ds, set) -> set.size());
watch.start();
DruidCoordinatorRuntimeParams balanceParams = tester.run(params);
DruidCoordinatorRuntimeParams assignParams = runner.run(params);
diff --git
a/server/src/test/java/org/apache/druid/server/coordinator/DruidCoordinatorTest.java
b/server/src/test/java/org/apache/druid/server/coordinator/DruidCoordinatorTest.java
index c2acef9e910..efa67668e8d 100644
---
a/server/src/test/java/org/apache/druid/server/coordinator/DruidCoordinatorTest.java
+++
b/server/src/test/java/org/apache/druid/server/coordinator/DruidCoordinatorTest.java
@@ -39,6 +39,7 @@ import org.apache.druid.java.util.emitter.core.Event;
import org.apache.druid.java.util.emitter.service.ServiceEmitter;
import org.apache.druid.java.util.emitter.service.ServiceMetricEvent;
import org.apache.druid.metadata.MetadataRuleManager;
+import org.apache.druid.metadata.MetadataRuleManagerConfig;
import org.apache.druid.metadata.SegmentsMetadataManager;
import org.apache.druid.metadata.segment.cache.NoopSegmentMetadataCache;
import org.apache.druid.rpc.indexing.OverlordClient;
@@ -68,6 +69,7 @@ import
org.apache.druid.server.coordinator.loading.TestLoadQueuePeon;
import
org.apache.druid.server.coordinator.rules.ForeverBroadcastDistributionRule;
import org.apache.druid.server.coordinator.rules.ForeverLoadRule;
import org.apache.druid.server.coordinator.rules.IntervalLoadRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.rules.Rule;
import org.apache.druid.server.coordinator.stats.Stats;
import org.apache.druid.server.http.BrokerDynamicConfigSyncer;
@@ -210,8 +212,8 @@ public class DruidCoordinatorTest
// Setup MetadataRuleManager
Rule foreverLoadRule = new ForeverLoadRule(ImmutableMap.of(tier, 2), null);
-
EasyMock.expect(metadataRuleManager.getRulesWithDefault(EasyMock.anyString()))
- .andReturn(ImmutableList.of(foreverLoadRule)).atLeastOnce();
+ EasyMock.expect(metadataRuleManager.getRulesSnapshot())
+ .andReturn(clusterDefaultRules(foreverLoadRule)).atLeastOnce();
metadataRuleManager.stop();
EasyMock.expectLastCall().once();
@@ -341,8 +343,8 @@ public class DruidCoordinatorTest
setupSegmentsMetadataMock(druidDataSources[0]);
-
EasyMock.expect(metadataRuleManager.getRulesWithDefault(EasyMock.anyString()))
- .andReturn(ImmutableList.of(hotTier, coldTier)).atLeastOnce();
+ EasyMock.expect(metadataRuleManager.getRulesSnapshot())
+ .andReturn(clusterDefaultRules(hotTier, coldTier)).atLeastOnce();
EasyMock.expect(serverInventoryView.getInventory())
.andReturn(ImmutableList.of(hotServer, coldServer))
@@ -424,8 +426,8 @@ public class DruidCoordinatorTest
setupSegmentsMetadataMock(druidDataSource);
final Rule broadcastDistributionRule = new
ForeverBroadcastDistributionRule();
-
EasyMock.expect(metadataRuleManager.getRulesWithDefault(EasyMock.anyString()))
-
.andReturn(ImmutableList.of(broadcastDistributionRule)).atLeastOnce();
+ EasyMock.expect(metadataRuleManager.getRulesSnapshot())
+
.andReturn(clusterDefaultRules(broadcastDistributionRule)).atLeastOnce();
EasyMock.expect(serverInventoryView.getInventory())
.andReturn(ImmutableList.of(hotServer, coldServer, brokerServer1,
brokerServer2, peonServer))
@@ -729,8 +731,8 @@ public class DruidCoordinatorTest
// Setup MetadataRuleManager
Rule intervalLoadRule = new
IntervalLoadRule(Intervals.of("2010-02-01/P1M"), ImmutableMap.of(hotTier, 1),
null);
Rule foreverLoadRule = new ForeverLoadRule(ImmutableMap.of(coldTier, 0),
null);
-
EasyMock.expect(metadataRuleManager.getRulesWithDefault(EasyMock.anyString()))
- .andReturn(ImmutableList.of(intervalLoadRule,
foreverLoadRule)).atLeastOnce();
+ EasyMock.expect(metadataRuleManager.getRulesSnapshot())
+ .andReturn(clusterDefaultRules(intervalLoadRule,
foreverLoadRule)).atLeastOnce();
metadataRuleManager.stop();
EasyMock.expectLastCall().once();
@@ -844,6 +846,18 @@ public class DruidCoordinatorTest
Assertions.assertEquals(Collections.emptyMap(),
result.getCompactionStates());
}
+ /**
+ * A rules snapshot in which the given rules are the cluster defaults, so
that they
+ * apply to every datasource.
+ */
+ private static RetentionRulesSnapshot clusterDefaultRules(Rule... rules)
+ {
+ return new RetentionRulesSnapshot(
+ Map.of(MetadataRuleManagerConfig.DEFAULT_RULE_NAME, List.of(rules)),
+ MetadataRuleManagerConfig.DEFAULT_RULE_NAME
+ );
+ }
+
private void setupSegmentsMetadataMock(DruidDataSource dataSource)
{
EasyMock.expect(segmentsMetadataManager.isPollingDatabasePeriodically()).andReturn(true).anyTimes();
diff --git
a/server/src/test/java/org/apache/druid/server/coordinator/duty/RunRulesPartialLoadPlacementTest.java
b/server/src/test/java/org/apache/druid/server/coordinator/duty/RunRulesPartialLoadPlacementTest.java
index 300f6cbecea..46e2284503e 100644
---
a/server/src/test/java/org/apache/druid/server/coordinator/duty/RunRulesPartialLoadPlacementTest.java
+++
b/server/src/test/java/org/apache/druid/server/coordinator/duty/RunRulesPartialLoadPlacementTest.java
@@ -26,6 +26,7 @@ import org.apache.druid.client.DruidServer;
import org.apache.druid.java.util.common.DateTimes;
import org.apache.druid.java.util.common.Intervals;
import org.apache.druid.java.util.common.concurrent.Execs;
+import org.apache.druid.metadata.MetadataRuleManagerConfig;
import org.apache.druid.segment.column.ColumnType;
import org.apache.druid.segment.column.RowSignature;
import org.apache.druid.server.coordination.ServerType;
@@ -40,6 +41,7 @@ import
org.apache.druid.server.coordinator.loading.TestLoadQueuePeon;
import org.apache.druid.server.coordinator.rules.CannotMatchBehavior;
import org.apache.druid.server.coordinator.rules.ForeverPartialLoadRule;
import org.apache.druid.server.coordinator.rules.PeriodPartialLoadRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.rules.Rule;
import
org.apache.druid.server.coordinator.rules.WildcardClusterGroupPartialLoadMatcher;
import org.apache.druid.server.coordinator.stats.CoordinatorRunStats;
@@ -206,11 +208,17 @@ public class RunRulesPartialLoadPlacementTest
private CoordinatorRunStats runRules(DruidCluster cluster, Rule rule,
DataSegment... segments)
{
final List<Rule> rules = Collections.singletonList(rule);
- final RunRules ruleRunner = new RunRules((ds, set) -> set.size(),
datasource -> rules);
+ final RunRules ruleRunner = new RunRules((ds, set) -> set.size());
DruidCoordinatorRuntimeParams params = DruidCoordinatorRuntimeParams
.builder()
.withDruidCluster(cluster)
+ .withRetentionRulesSnapshot(
+ new RetentionRulesSnapshot(
+ Map.of(MetadataRuleManagerConfig.DEFAULT_RULE_NAME, rules),
+ MetadataRuleManagerConfig.DEFAULT_RULE_NAME
+ )
+ )
.withUsedSegments(segments)
.withBalancerStrategy(balancerStrategy)
.withDynamicConfigs(
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 6070c3d89f6..8cda9f8546a 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
@@ -23,6 +23,7 @@ import com.google.common.collect.ImmutableMap;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.MoreExecutors;
+import org.apache.druid.audit.AuditInfo;
import org.apache.druid.client.DruidServer;
import org.apache.druid.java.util.common.DateTimes;
import org.apache.druid.java.util.common.Intervals;
@@ -32,7 +33,7 @@ import org.apache.druid.java.util.emitter.EmittingLogger;
import org.apache.druid.java.util.emitter.core.EventMap;
import org.apache.druid.java.util.emitter.service.AlertEvent;
import org.apache.druid.java.util.metrics.StubServiceEmitter;
-import org.apache.druid.metadata.MetadataRuleManager;
+import org.apache.druid.metadata.MetadataRuleManagerConfig;
import org.apache.druid.segment.IndexIO;
import org.apache.druid.server.coordination.ServerType;
import org.apache.druid.server.coordinator.CoordinatorDynamicConfig;
@@ -50,6 +51,8 @@ import
org.apache.druid.server.coordinator.loading.TestLoadQueuePeon;
import org.apache.druid.server.coordinator.rules.ForeverLoadRule;
import org.apache.druid.server.coordinator.rules.IntervalDropRule;
import org.apache.druid.server.coordinator.rules.IntervalLoadRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
+import org.apache.druid.server.coordinator.rules.Rule;
import org.apache.druid.server.coordinator.stats.CoordinatorRunStats;
import org.apache.druid.server.coordinator.stats.Dimension;
import org.apache.druid.server.coordinator.stats.RowKey;
@@ -69,6 +72,7 @@ import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
+import java.util.Map;
import java.util.Set;
import java.util.concurrent.TimeUnit;
@@ -80,11 +84,12 @@ public class RunRulesTest
private static final long SERVER_SIZE_10GB = 10L << 30;
private static final String DATASOURCE = "test";
private static final RowKey DATASOURCE_STAT_KEY =
RowKey.of(Dimension.DATASOURCE, DATASOURCE);
+ private static final AuditInfo AUDIT_INFO = new AuditInfo("test", "id",
"test", "127.0.0.1");
private LoadQueuePeon mockPeon;
private RunRules ruleRunner;
private StubServiceEmitter emitter;
- private MetadataRuleManager databaseRuleManager;
+ private RetentionRulesSnapshot rulesSnapshot;
private SegmentLoadQueueManager loadQueueManager;
private final List<DataSegment> usedSegments =
CreateDataSegments.ofDatasource(DATASOURCE)
@@ -101,8 +106,8 @@ public class RunRulesTest
mockPeon = EasyMock.createMock(LoadQueuePeon.class);
emitter = StubServiceEmitter.createStarted();
EmittingLogger.registerEmitter(emitter);
- databaseRuleManager = EasyMock.createMock(MetadataRuleManager.class);
- ruleRunner = new RunRules((ds, set) -> set.size(),
databaseRuleManager::getRulesWithDefault);
+ rulesSnapshot = new RetentionRulesSnapshot(Map.of(),
MetadataRuleManagerConfig.DEFAULT_RULE_NAME);
+ ruleRunner = new RunRules((ds, set) -> set.size());
loadQueueManager = new SegmentLoadQueueManager(null, null);
balancerExecutor = MoreExecutors.listeningDecorator(Execs.multiThreaded(1,
"RunRulesTest-%d"));
}
@@ -111,7 +116,19 @@ public class RunRulesTest
public void tearDown()
{
balancerExecutor.shutdown();
- EasyMock.verify(databaseRuleManager);
+ }
+
+ /**
+ * Sets the rules that {@link #createCoordinatorRuntimeParams} will snapshot
into the
+ * params for this test. They are set as cluster defaults so that they apply
to every
+ * datasource.
+ */
+ private void setRetentionRules(Rule... rules)
+ {
+ rulesSnapshot = new RetentionRulesSnapshot(
+ Map.of(MetadataRuleManagerConfig.DEFAULT_RULE_NAME, List.of(rules)),
+ MetadataRuleManagerConfig.DEFAULT_RULE_NAME
+ );
}
/**
@@ -127,15 +144,9 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Collections.singletonList(
- new IntervalLoadRule(
- Intervals.of("2012-01-01/2012-01-02"),
- ImmutableMap.of("normal", 2),
- null
- )
- )).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(
+ new IntervalLoadRule(Intervals.of("2012-01-01/2012-01-02"),
ImmutableMap.of("normal", 2), null)
+ );
// server1 has all the segments already loaded
final DruidServer server1 = createHistorical("server1", "normal");
@@ -184,15 +195,13 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Collections.singletonList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"),
- ImmutableMap.of("hot", 2, "normal", 2),
- null
- )
- )).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"),
+ ImmutableMap.of("hot", 2, "normal", 2),
+ null
+ )
+ );
final DruidServer serverHot1 = createHistorical("serverHot", "hot");
final DruidServer serverHot2 = createHistorical("serverHot2", "hot");
@@ -247,25 +256,23 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Lists.newArrayList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T06:00:00.000Z"),
- ImmutableMap.of("hot", 1),
- null
- ),
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
- ImmutableMap.of("normal", 1),
- null
- ),
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"),
- ImmutableMap.of("cold", 1),
- null
- )
- )).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T06:00:00.000Z"),
+ ImmutableMap.of("hot", 1),
+ null
+ ),
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
+ ImmutableMap.of("normal", 1),
+ null
+ ),
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"),
+ ImmutableMap.of("cold", 1),
+ null
+ )
+ );
DruidCluster druidCluster = DruidCluster
.builder()
@@ -339,6 +346,7 @@ public class RunRulesTest
return DruidCoordinatorRuntimeParams
.builder()
.withDruidCluster(druidCluster)
+ .withRetentionRulesSnapshot(rulesSnapshot)
.withUsedSegments(dataSegments);
}
@@ -354,21 +362,18 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Lists.newArrayList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T06:00:00.000Z"),
- ImmutableMap.of("hot", 2),
- null
- ),
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"),
- ImmutableMap.of("cold", 1),
- null
- )
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T06:00:00.000Z"),
+ ImmutableMap.of("hot", 2),
+ null
+ ),
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"),
+ ImmutableMap.of("cold", 1),
+ null
)
- ).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ );
DruidCluster druidCluster = DruidCluster
.builder()
@@ -402,21 +407,18 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Lists.newArrayList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
- ImmutableMap.of("hot", 1),
- null
- ),
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"),
- ImmutableMap.of("normal", 1),
- null
- )
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
+ ImmutableMap.of("hot", 1),
+ null
+ ),
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"),
+ ImmutableMap.of("normal", 1),
+ null
)
- ).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ );
DruidServer normServer = createHistorical("serverNorm", "normal");
for (DataSegment segment : usedSegments) {
@@ -450,21 +452,18 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Lists.newArrayList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
- ImmutableMap.of("hot", 1),
- null
- ),
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"),
- ImmutableMap.of("normal", 1),
- null
- )
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
+ ImmutableMap.of("hot", 1),
+ null
+ ),
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"),
+ ImmutableMap.of("normal", 1),
+ null
)
- ).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ );
DruidCluster druidCluster = DruidCluster
.builder()
@@ -485,22 +484,17 @@ public class RunRulesTest
public void testRunRuleDoesNotExist()
{
- EasyMock
- .expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject()))
- .andReturn(
- Collections.singletonList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-02T00:00:00.000Z/2012-01-03T00:00:00.000Z"),
- ImmutableMap.of("normal", 1),
- null
- )
- )
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-02T00:00:00.000Z/2012-01-03T00:00:00.000Z"),
+ ImmutableMap.of("normal", 1),
+ null
)
- .atLeastOnce();
+ );
EasyMock.expect(mockPeon.getSegmentsInQueue()).andReturn(Collections.emptySet()).anyTimes();
EasyMock.expect(mockPeon.getSegmentsMarkedToDrop()).andReturn(Collections.emptySet()).anyTimes();
- EasyMock.replay(databaseRuleManager, mockPeon);
+ EasyMock.replay(mockPeon);
DruidCluster druidCluster = DruidCluster
.builder()
@@ -533,17 +527,14 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Lists.newArrayList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
- ImmutableMap.of("normal", 1),
- null
- ),
- new
IntervalDropRule(Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"))
- )
- ).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
+ ImmutableMap.of("normal", 1),
+ null
+ ),
+ new
IntervalDropRule(Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"))
+ );
DruidServer server = createHistorical("serverNorm", "normal");
for (DataSegment segment : usedSegments) {
@@ -571,17 +562,14 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Lists.newArrayList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
- ImmutableMap.of("normal", 1),
- null
- ),
- new
IntervalDropRule(Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"))
- )
- ).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
+ ImmutableMap.of("normal", 1),
+ null
+ ),
+ new
IntervalDropRule(Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"))
+ );
DruidServer server1 = createHistorical("serverNorm", "normal");
server1.addDataSegment(usedSegments.get(0));
@@ -628,17 +616,14 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Lists.newArrayList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
- ImmutableMap.of("hot", 1),
- null
- ),
- new
IntervalDropRule(Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"))
- )
- ).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
+ ImmutableMap.of("hot", 1),
+ null
+ ),
+ new
IntervalDropRule(Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"))
+ );
DruidServer server1 = createHistorical("server1", "hot");
server1.addDataSegment(usedSegments.get(0));
@@ -672,17 +657,14 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Lists.newArrayList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
- ImmutableMap.of("hot", 1),
- null
- ),
- new
IntervalDropRule(Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"))
- )
- ).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T12:00:00.000Z"),
+ ImmutableMap.of("hot", 1),
+ null
+ ),
+ new
IntervalDropRule(Intervals.of("2012-01-01T00:00:00.000Z/2012-01-02T00:00:00.000Z"))
+ );
DruidServer server1 = createHistorical("server1", "hot");
DruidServer server2 = createHistorical("serverNorm2", "normal");
@@ -711,19 +693,13 @@ public class RunRulesTest
@Test
public void testDropServerActuallyServesSegment()
{
- EasyMock
- .expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject()))
- .andReturn(
- Collections.singletonList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T01:00:00.000Z"),
- ImmutableMap.of("normal", 0),
- null
- )
- )
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2012-01-01T01:00:00.000Z"),
+ ImmutableMap.of("normal", 0),
+ null
)
- .atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ );
DruidServer server1 = createHistorical("server1", "normal");
server1.addDataSegment(usedSegments.get(0));
@@ -777,19 +753,13 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
- EasyMock
- .expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject()))
- .andReturn(
- Collections.singletonList(
- new IntervalLoadRule(
-
Intervals.of("2012-01-01T00:00:00.000Z/2013-01-01T00:00:00.000Z"),
- ImmutableMap.of("hot", 2),
- null
- )
- )
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01T00:00:00.000Z/2013-01-01T00:00:00.000Z"),
+ ImmutableMap.of("hot", 2),
+ null
)
- .atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ );
DruidCluster druidCluster = DruidCluster
.builder()
@@ -854,19 +824,13 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
- EasyMock
- .expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject()))
- .andReturn(
- Collections.singletonList(
- new IntervalLoadRule(
- Intervals.of("2012-01-01/2013-01-01"),
- ImmutableMap.of("hot", 1, DruidServer.DEFAULT_TIER, 1),
- null
- )
- )
+ setRetentionRules(
+ new IntervalLoadRule(
+ Intervals.of("2012-01-01/2013-01-01"),
+ ImmutableMap.of("hot", 1, DruidServer.DEFAULT_TIER, 1),
+ null
)
- .atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ );
final DruidCluster druidCluster = DruidCluster
.builder()
@@ -907,19 +871,9 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
- EasyMock
- .expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject()))
- .andReturn(
- Collections.singletonList(
- new IntervalLoadRule(
- Intervals.of("2012-01-01/2013-01-02"),
- ImmutableMap.of("normal", 1),
- null
- )
- )
- )
- .atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(
+ new IntervalLoadRule(Intervals.of("2012-01-01/2013-01-02"),
ImmutableMap.of("normal", 1), null)
+ );
DataSegment overFlowSegment = new DataSegment(
"test",
@@ -996,9 +950,7 @@ public class RunRulesTest
EasyMock.expectLastCall().once();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Collections.singletonList(new
ForeverLoadRule(ImmutableMap.of(DruidServer.DEFAULT_TIER, 1),
null))).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(new
ForeverLoadRule(ImmutableMap.of(DruidServer.DEFAULT_TIER, 1), null));
DruidCluster druidCluster = DruidCluster.builder().add(
createServerHolder("serverHot", DruidServer.DEFAULT_TIER, mockPeon)
@@ -1033,11 +985,7 @@ public class RunRulesTest
{
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Collections.singletonList(
- new ForeverLoadRule(ImmutableMap.of(DruidServer.DEFAULT_TIER, 3),
null)
- )).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(new
ForeverLoadRule(ImmutableMap.of(DruidServer.DEFAULT_TIER, 3), null));
DataSegment dataSegment = new DataSegment(
"test",
@@ -1096,11 +1044,7 @@ public class RunRulesTest
EasyMock.expectLastCall().atLeastOnce();
mockEmptyPeon();
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Collections.singletonList(
- new ForeverLoadRule(ImmutableMap.of(DruidServer.DEFAULT_TIER, 1),
null)
- )).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(new
ForeverLoadRule(ImmutableMap.of(DruidServer.DEFAULT_TIER, 1), null));
DataSegment dataSegment = new DataSegment(
"test",
@@ -1148,14 +1092,7 @@ public class RunRulesTest
{
mockEmptyPeon();
int numReplicants = 1;
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Collections.singletonList(
- new ForeverLoadRule(
- ImmutableMap.of(DruidServer.DEFAULT_TIER, numReplicants),
- null
- )
- )).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(new
ForeverLoadRule(ImmutableMap.of(DruidServer.DEFAULT_TIER, numReplicants),
null));
final DataSegment dataSegment = new DataSegment(
"test",
@@ -1210,14 +1147,7 @@ public class RunRulesTest
{
mockEmptyPeon();
int numReplicants = 1;
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Collections.singletonList(
- new ForeverLoadRule(
- ImmutableMap.of(DruidServer.DEFAULT_TIER, numReplicants),
- null
- )
- )).atLeastOnce();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(new
ForeverLoadRule(ImmutableMap.of(DruidServer.DEFAULT_TIER, numReplicants),
null));
DataSegment dataSegment = new DataSegment(
"test",
@@ -1262,12 +1192,7 @@ public class RunRulesTest
@Test
public void testSegmentWithZeroRequiredReplicasHasZeroReplicationFactor()
{
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(EasyMock.anyObject())).andReturn(
- Collections.singletonList(
- new ForeverLoadRule(Collections.emptyMap(), false)
- )
- ).anyTimes();
- EasyMock.replay(databaseRuleManager);
+ setRetentionRules(new ForeverLoadRule(Collections.emptyMap(), false));
final DruidCluster cluster = DruidCluster
.builder()
diff --git
a/server/src/test/java/org/apache/druid/server/coordinator/duty/UnloadUnusedSegmentsTest.java
b/server/src/test/java/org/apache/druid/server/coordinator/duty/UnloadUnusedSegmentsTest.java
index 7e7f9b6bb27..e468c8537e5 100644
---
a/server/src/test/java/org/apache/druid/server/coordinator/duty/UnloadUnusedSegmentsTest.java
+++
b/server/src/test/java/org/apache/druid/server/coordinator/duty/UnloadUnusedSegmentsTest.java
@@ -27,7 +27,7 @@ import org.apache.druid.client.ImmutableDruidDataSource;
import org.apache.druid.client.ImmutableDruidServer;
import org.apache.druid.client.ImmutableDruidServerTests;
import org.apache.druid.java.util.common.DateTimes;
-import org.apache.druid.metadata.MetadataRuleManager;
+import org.apache.druid.metadata.MetadataRuleManagerConfig;
import org.apache.druid.server.coordination.ServerType;
import org.apache.druid.server.coordinator.DruidCluster;
import org.apache.druid.server.coordinator.DruidCoordinator;
@@ -37,6 +37,8 @@ import
org.apache.druid.server.coordinator.loading.SegmentLoadQueueManager;
import org.apache.druid.server.coordinator.loading.TestLoadQueuePeon;
import
org.apache.druid.server.coordinator.rules.ForeverBroadcastDistributionRule;
import org.apache.druid.server.coordinator.rules.ForeverLoadRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
+import org.apache.druid.server.coordinator.rules.Rule;
import org.apache.druid.server.coordinator.stats.CoordinatorRunStats;
import org.apache.druid.server.coordinator.stats.Stats;
import org.apache.druid.timeline.DataSegment;
@@ -52,6 +54,7 @@ import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
+import java.util.Map;
import java.util.Set;
public class UnloadUnusedSegmentsTest
@@ -72,7 +75,6 @@ public class UnloadUnusedSegmentsTest
private List<ImmutableDruidDataSource> dataSources;
private List<ImmutableDruidDataSource> dataSourcesForRealtime;
private final String broadcastDatasource = "broadcastDatasource";
- private MetadataRuleManager databaseRuleManager;
private SegmentLoadQueueManager loadQueueManager;
@BeforeEach
@@ -83,7 +85,6 @@ public class UnloadUnusedSegmentsTest
historicalServerTier2 = EasyMock.createMock(ImmutableDruidServer.class);
brokerServer = EasyMock.createMock(ImmutableDruidServer.class);
indexerServer = EasyMock.createMock(ImmutableDruidServer.class);
- databaseRuleManager = EasyMock.createMock(MetadataRuleManager.class);
loadQueueManager = new SegmentLoadQueueManager(null, null);
DateTime start1 = DateTimes.of("2012-01-01");
@@ -185,11 +186,9 @@ public class UnloadUnusedSegmentsTest
EasyMock.verify(historicalServerTier2);
EasyMock.verify(brokerServer);
EasyMock.verify(indexerServer);
- EasyMock.verify(databaseRuleManager);
}
- @Test
- public void test_unloadUnusedSegmentsFromAllServers()
+ private void mockAllServers()
{
mockDruidServer(
historicalServer,
@@ -234,8 +233,29 @@ public class UnloadUnusedSegmentsTest
// Mock stuff that the coordinator needs
mockCoordinator(coordinator);
+ }
+
+ private DruidCluster buildCluster()
+ {
+ return DruidCluster
+ .builder()
+ .addTier(
+ DruidServer.DEFAULT_TIER,
+ new ServerHolder(historicalServer, historicalPeon, false)
+ )
+ .addTier(
+ "tier2",
+ new ServerHolder(historicalServerTier2, historicalTier2Peon, false)
+ )
+ .addBrokers(new ServerHolder(brokerServer, brokerPeon, false))
+ .addRealtimes(new ServerHolder(indexerServer, indexerPeon, false))
+ .build();
+ }
- mockRuleManager(databaseRuleManager);
+ @Test
+ public void test_unloadUnusedSegmentsFromAllServers()
+ {
+ mockAllServers();
// We keep datasource2 segments only, drop datasource1 and
broadcastDatasource from all servers
// realtimeSegment is intentionally missing from the set, to match how a
realtime tasks's unpublished segments
@@ -244,26 +264,13 @@ public class UnloadUnusedSegmentsTest
DruidCoordinatorRuntimeParams params = DruidCoordinatorRuntimeParams
.builder()
- .withDruidCluster(
- DruidCluster
- .builder()
- .addTier(
- DruidServer.DEFAULT_TIER,
- new ServerHolder(historicalServer, historicalPeon, false)
- )
- .addTier(
- "tier2",
- new ServerHolder(historicalServerTier2,
historicalTier2Peon, false)
- )
- .addBrokers(new ServerHolder(brokerServer, brokerPeon, false))
- .addRealtimes(new ServerHolder(indexerServer, indexerPeon,
false))
- .build()
- )
+ .withDruidCluster(buildCluster())
.withUsedSegments(usedSegments)
+ .withRetentionRulesSnapshot(rulesSnapshot())
.withBroadcastDatasources(Collections.singleton(broadcastDatasource))
.build();
- params = new UnloadUnusedSegments(loadQueueManager,
databaseRuleManager::getRulesWithDefault).run(params);
+ params = new UnloadUnusedSegments(loadQueueManager).run(params);
CoordinatorRunStats stats = params.getCoordinatorStats();
// We drop segment1 and broadcast1 from all servers, realtimeSegment is
not dropped by the indexer
@@ -274,6 +281,39 @@ public class UnloadUnusedSegmentsTest
Assertions.assertEquals(1L, stats.getSegmentStat(Stats.Segments.UNNEEDED,
"tier2", broadcastDatasource));
}
+ @Test
+ public void test_broadcastSegmentsAreDropped_ifSnapshotHasNoBroadcastRule()
+ {
+ mockAllServers();
+
+ final Set<DataSegment> usedSegments = ImmutableSet.of(segment2);
+
+ // The snapshot loads broadcastDatasource like any other datasource rather
than
+ // broadcasting it, so the realtime server must be skipped for it
+ final RetentionRulesSnapshot rulesWithoutBroadcast = new
RetentionRulesSnapshot(
+ Map.of(
+ MetadataRuleManagerConfig.DEFAULT_RULE_NAME,
+ List.of(new
ForeverLoadRule(ImmutableMap.of(DruidServer.DEFAULT_TIER, 1, "tier2", 1), null))
+ ),
+ MetadataRuleManagerConfig.DEFAULT_RULE_NAME
+ );
+
+ DruidCoordinatorRuntimeParams params = DruidCoordinatorRuntimeParams
+ .builder()
+ .withDruidCluster(buildCluster())
+ .withUsedSegments(usedSegments)
+ .withRetentionRulesSnapshot(rulesWithoutBroadcast)
+ .build();
+
+ params = new UnloadUnusedSegments(loadQueueManager).run(params);
+ final CoordinatorRunStats stats = params.getCoordinatorStats();
+
+ // Two broadcast segments are dropped from the historical, and none from
the realtime
+ // server, one fewer than when the snapshot marks the datasource as
broadcast
+ Assertions.assertEquals(2L, stats.getSegmentStat(Stats.Segments.UNNEEDED,
DruidServer.DEFAULT_TIER, broadcastDatasource));
+ Assertions.assertEquals(1L, stats.getSegmentStat(Stats.Segments.UNNEEDED,
"tier2", broadcastDatasource));
+ }
+
private static void mockDruidServer(
ImmutableDruidServer druidServer,
ServerType serverType,
@@ -307,35 +347,19 @@ public class UnloadUnusedSegmentsTest
EasyMock.replay(coordinator);
}
- private static void mockRuleManager(MetadataRuleManager metadataRuleManager)
+ private static RetentionRulesSnapshot rulesSnapshot()
{
-
EasyMock.expect(metadataRuleManager.getRulesWithDefault("datasource1")).andReturn(
- Collections.singletonList(
- new ForeverLoadRule(
- ImmutableMap.of(
- DruidServer.DEFAULT_TIER, 1,
- "tier2", 1
- ),
- null
- )
- )).anyTimes();
-
-
EasyMock.expect(metadataRuleManager.getRulesWithDefault("datasource2")).andReturn(
- Collections.singletonList(
- new ForeverLoadRule(
- ImmutableMap.of(
- DruidServer.DEFAULT_TIER, 1,
- "tier2", 1
- ),
- null
- )
- )).anyTimes();
-
-
EasyMock.expect(metadataRuleManager.getRulesWithDefault("broadcastDatasource")).andReturn(
- Collections.singletonList(
- new ForeverBroadcastDistributionRule()
- )).anyTimes();
-
- EasyMock.replay(metadataRuleManager);
+ final List<Rule> loadOnBothTiers = Collections.singletonList(
+ new ForeverLoadRule(ImmutableMap.of(DruidServer.DEFAULT_TIER, 1,
"tier2", 1), null)
+ );
+ return new RetentionRulesSnapshot(
+ Map.of(
+ "datasource1", loadOnBothTiers,
+ "datasource2", loadOnBothTiers,
+ "broadcastDatasource", Collections.singletonList(new
ForeverBroadcastDistributionRule())
+ ),
+ // Has no entry above, so every datasource is governed solely by its
own rules
+ "unusedDefaultDatasource"
+ );
}
}
diff --git
a/server/src/test/java/org/apache/druid/server/coordinator/rules/LoadRuleTest.java
b/server/src/test/java/org/apache/druid/server/coordinator/rules/LoadRuleTest.java
index 3022bd6cbe7..344c140d844 100644
---
a/server/src/test/java/org/apache/druid/server/coordinator/rules/LoadRuleTest.java
+++
b/server/src/test/java/org/apache/druid/server/coordinator/rules/LoadRuleTest.java
@@ -669,6 +669,48 @@ public class LoadRuleTest
Assertions.assertEquals(2L, stats.getSegmentStat(Stats.Segments.ASSIGNED,
Tier.T2, TestDataSource.WIKI));
}
+ @Test
+ public void test_allReplicasAreDropped_ifAliasTierIsUnresolved()
+ {
+ final DataSegment segment = createDataSegment(TestDataSource.WIKI);
+
+ DruidCluster cluster = DruidCluster
+ .builder()
+ .addTier(Tier.T2, createServer(Tier.T2, segment),
createServer(Tier.T2, segment))
+ .build();
+
+ // Aliases are configured, but none of them defines T1, so it resolves to
no real tier
+ LoadRule rule = loadForever(ImmutableMap.of(Tier.T1, 1));
+ CoordinatorRunStats stats = runRuleAndGetStats(
+ rule,
+ segment,
+ makeCoordinatorRuntimeParams(cluster, ImmutableMap.of(Tier.T3,
Set.of(Tier.T2)), segment)
+ );
+
+ Assertions.assertEquals(2L, stats.getSegmentStat(Stats.Segments.DROPPED,
Tier.T2, TestDataSource.WIKI));
+ Assertions.assertEquals(0L, stats.getSegmentStat(Stats.Segments.ASSIGNED,
Tier.T1, TestDataSource.WIKI));
+ }
+
+ @Test
+ public void test_onlySurplusIsDropped_ifAliasTierIsResolved()
+ {
+ final DataSegment segment = createDataSegment(TestDataSource.WIKI);
+
+ DruidCluster cluster = DruidCluster
+ .builder()
+ .addTier(Tier.T2, createServer(Tier.T2, segment),
createServer(Tier.T2, segment))
+ .build();
+
+ LoadRule rule = loadForever(ImmutableMap.of(Tier.T1, 1));
+ CoordinatorRunStats stats = runRuleAndGetStats(
+ rule,
+ segment,
+ makeCoordinatorRuntimeParams(cluster, ImmutableMap.of(Tier.T1,
Set.of(Tier.T2)), segment)
+ );
+
+ Assertions.assertEquals(1L, stats.getSegmentStat(Stats.Segments.DROPPED,
Tier.T2, TestDataSource.WIKI));
+ }
+
private DruidCoordinatorRuntimeParams makeCoordinatorRuntimeParams(
DruidCluster druidCluster,
Map<String, Set<String>> historicalTierAliases,
diff --git
a/server/src/test/java/org/apache/druid/server/coordinator/rules/RetentionRulesSnapshotTest.java
b/server/src/test/java/org/apache/druid/server/coordinator/rules/RetentionRulesSnapshotTest.java
new file mode 100644
index 00000000000..ded6d022c18
--- /dev/null
+++
b/server/src/test/java/org/apache/druid/server/coordinator/rules/RetentionRulesSnapshotTest.java
@@ -0,0 +1,168 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.druid.server.coordinator.rules;
+
+import org.apache.druid.client.DruidServer;
+import org.apache.druid.java.util.common.Intervals;
+import org.apache.druid.segment.TestDataSource;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+class RetentionRulesSnapshotTest
+{
+ private static final Rule DATASOURCE_RULE = new IntervalLoadRule(
+ Intervals.of("2012-01-01/2012-01-02"),
+ Map.of(DruidServer.DEFAULT_TIER, 2),
+ null
+ );
+ private static final Rule DEFAULT_RULE = new
ForeverLoadRule(Map.of(DruidServer.DEFAULT_TIER, 1), null);
+ /**
+ * Deliberately not {@code _default}, so that these tests fail if the
snapshot ever assumes
+ * the conventional name instead of honouring the configured one.
+ */
+ private static final String DEFAULT_DATASOURCE = "configuredDefaultRules";
+
+ @Test
+ void testDatasourceRulesArePrependedToClusterDefaults()
+ {
+ final RetentionRulesSnapshot rules = new RetentionRulesSnapshot(
+ Map.of(
+ TestDataSource.WIKI, List.of(DATASOURCE_RULE),
+ DEFAULT_DATASOURCE, List.of(DEFAULT_RULE)
+ ),
+ DEFAULT_DATASOURCE
+ );
+
+ Assertions.assertEquals(
+ List.of(DATASOURCE_RULE, DEFAULT_RULE),
+ rules.getEffectiveRules(TestDataSource.WIKI)
+ );
+ }
+
+ @Test
+ void testDatasourceWithNoRulesResolvesToClusterDefaults()
+ {
+ final RetentionRulesSnapshot rules = new RetentionRulesSnapshot(
+ Map.of(
+ TestDataSource.WIKI, List.of(DATASOURCE_RULE),
+ DEFAULT_DATASOURCE, List.of(DEFAULT_RULE)
+ ),
+ DEFAULT_DATASOURCE
+ );
+
+ Assertions.assertEquals(List.of(DEFAULT_RULE),
rules.getEffectiveRules(TestDataSource.KOALA));
+ }
+
+ @Test
+ void testSnapshotIsUnaffectedByLaterChangesToSourceMap()
+ {
+ final Map<String, List<Rule>> source = new HashMap<>();
+ source.put(TestDataSource.WIKI, List.of(DATASOURCE_RULE));
+ source.put(DEFAULT_DATASOURCE, List.of(DEFAULT_RULE));
+
+ final RetentionRulesSnapshot rules = new RetentionRulesSnapshot(source,
DEFAULT_DATASOURCE);
+
+ source.put(TestDataSource.WIKI, List.of(DEFAULT_RULE));
+ source.put(TestDataSource.KOALA, List.of(DATASOURCE_RULE));
+
+ Assertions.assertEquals(
+ List.of(DATASOURCE_RULE, DEFAULT_RULE),
+ rules.getEffectiveRules(TestDataSource.WIKI)
+ );
+ Assertions.assertEquals(List.of(DEFAULT_RULE),
rules.getEffectiveRules(TestDataSource.KOALA));
+ }
+
+ @Test
+ void testSnapshotIsUnaffectedByLaterChangesToSourceRuleList()
+ {
+ final List<Rule> wikiRules = new ArrayList<>(List.of(DATASOURCE_RULE));
+ final RetentionRulesSnapshot rules = new RetentionRulesSnapshot(
+ Map.of(
+ TestDataSource.WIKI, wikiRules,
+ DEFAULT_DATASOURCE, List.of(DEFAULT_RULE)
+ ),
+ DEFAULT_DATASOURCE
+ );
+
+ wikiRules.clear();
+
+ Assertions.assertEquals(List.of(DATASOURCE_RULE),
rules.getOverrideRules(TestDataSource.WIKI));
+ Assertions.assertEquals(
+ List.of(DATASOURCE_RULE, DEFAULT_RULE),
+ rules.getEffectiveRules(TestDataSource.WIKI)
+ );
+ }
+
+ @Test
+ void testClusterDefaultsAreEmptyWhenDefaultDatasourceHasNoRules()
+ {
+ // The default datasource has no rules until
SQLMetadataRuleManager.createDefaultRule() has run
+ final RetentionRulesSnapshot rules = new RetentionRulesSnapshot(
+ Map.of(TestDataSource.WIKI, List.of(DATASOURCE_RULE)),
+ DEFAULT_DATASOURCE
+ );
+
+ Assertions.assertEquals(List.of(DATASOURCE_RULE),
rules.getEffectiveRules(TestDataSource.WIKI));
+ Assertions.assertEquals(List.of(),
rules.getEffectiveRules(TestDataSource.KOALA));
+ }
+
+ @Test
+ void testDefaultDatasourceRulesAreNotAppendedToThemselves()
+ {
+ final RetentionRulesSnapshot rules = new RetentionRulesSnapshot(
+ Map.of(
+ TestDataSource.WIKI, List.of(DATASOURCE_RULE),
+ DEFAULT_DATASOURCE, List.of(DEFAULT_RULE)
+ ),
+ DEFAULT_DATASOURCE
+ );
+
+ Assertions.assertEquals(List.of(DEFAULT_RULE),
rules.getEffectiveRules(DEFAULT_DATASOURCE));
+ }
+
+ @Test
+ void testClusterDefaultsApplyToEveryDatasourceWithNoRulesOfItsOwn()
+ {
+ final RetentionRulesSnapshot rules = new RetentionRulesSnapshot(
+ Map.of(DEFAULT_DATASOURCE, List.of(DEFAULT_RULE)),
+ DEFAULT_DATASOURCE
+ );
+
+ Assertions.assertEquals(List.of(DEFAULT_RULE),
rules.getEffectiveRules(TestDataSource.WIKI));
+ Assertions.assertEquals(List.of(DEFAULT_RULE),
rules.getEffectiveRules(TestDataSource.KOALA));
+
+ // No datasource has override rules of its own
+ Assertions.assertEquals(List.of(),
rules.getOverrideRules(TestDataSource.WIKI));
+ }
+
+ @Test
+ void testEmptySnapshotHasNoRulesForAnyDatasource()
+ {
+ final RetentionRulesSnapshot rules = new RetentionRulesSnapshot(Map.of(),
DEFAULT_DATASOURCE);
+
+ Assertions.assertEquals(List.of(),
rules.getEffectiveRules(TestDataSource.WIKI));
+ Assertions.assertEquals(Map.of(), rules.getAllRules());
+ }
+}
diff --git
a/server/src/test/java/org/apache/druid/server/coordinator/simulate/TestMetadataRuleManager.java
b/server/src/test/java/org/apache/druid/server/coordinator/simulate/TestMetadataRuleManager.java
index 9bb877112cf..2e50923bb8a 100644
---
a/server/src/test/java/org/apache/druid/server/coordinator/simulate/TestMetadataRuleManager.java
+++
b/server/src/test/java/org/apache/druid/server/coordinator/simulate/TestMetadataRuleManager.java
@@ -21,10 +21,11 @@ package org.apache.druid.server.coordinator.simulate;
import org.apache.druid.audit.AuditInfo;
import org.apache.druid.metadata.MetadataRuleManager;
+import org.apache.druid.metadata.MetadataRuleManagerConfig;
import org.apache.druid.server.coordinator.rules.ForeverLoadRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.rules.Rule;
-import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
@@ -34,7 +35,7 @@ public class TestMetadataRuleManager implements
MetadataRuleManager
{
private final Map<String, List<Rule>> rules = new HashMap<>();
- private static final String DEFAULT_DATASOURCE = "_default";
+ private static final String DEFAULT_DATASOURCE =
MetadataRuleManagerConfig.DEFAULT_RULE_NAME;
public TestMetadataRuleManager()
{
@@ -63,30 +64,9 @@ public class TestMetadataRuleManager implements
MetadataRuleManager
}
@Override
- public Map<String, List<Rule>> getAllRules()
+ public RetentionRulesSnapshot getRulesSnapshot()
{
- return rules;
- }
-
- @Override
- public List<Rule> getRules(final String dataSource)
- {
- List<Rule> retVal = rules.get(dataSource);
- return retVal == null ? new ArrayList<>() : retVal;
- }
-
- @Override
- public List<Rule> getRulesWithDefault(final String dataSource)
- {
- List<Rule> retVal = new ArrayList<>();
- final Map<String, List<Rule>> theRules = rules;
- if (theRules.get(dataSource) != null) {
- retVal.addAll(theRules.get(dataSource));
- }
- if (theRules.get(DEFAULT_DATASOURCE) != null) {
- retVal.addAll(theRules.get(DEFAULT_DATASOURCE));
- }
- return retVal;
+ return new RetentionRulesSnapshot(rules, DEFAULT_DATASOURCE);
}
@Override
diff --git
a/server/src/test/java/org/apache/druid/server/http/DataSourcesResourceTest.java
b/server/src/test/java/org/apache/druid/server/http/DataSourcesResourceTest.java
index 184c29ccd65..b525c4f126a 100644
---
a/server/src/test/java/org/apache/druid/server/http/DataSourcesResourceTest.java
+++
b/server/src/test/java/org/apache/druid/server/http/DataSourcesResourceTest.java
@@ -44,6 +44,7 @@ import org.apache.druid.java.util.common.DateTimes;
import org.apache.druid.java.util.common.Intervals;
import
org.apache.druid.java.util.http.client.response.StringFullResponseHolder;
import org.apache.druid.metadata.MetadataRuleManager;
+import org.apache.druid.metadata.MetadataRuleManagerConfig;
import org.apache.druid.metadata.SegmentsMetadataManager;
import org.apache.druid.query.SegmentDescriptor;
import org.apache.druid.query.TableDataSource;
@@ -59,6 +60,7 @@ import
org.apache.druid.server.coordinator.rules.ExactProjectionPartialLoadMatch
import org.apache.druid.server.coordinator.rules.IntervalDropRule;
import org.apache.druid.server.coordinator.rules.IntervalLoadRule;
import org.apache.druid.server.coordinator.rules.IntervalPartialLoadRule;
+import org.apache.druid.server.coordinator.rules.RetentionRulesSnapshot;
import org.apache.druid.server.coordinator.rules.Rule;
import
org.apache.druid.server.coordinator.rules.WildcardClusterGroupPartialLoadMatcher;
import org.apache.druid.server.security.Access;
@@ -672,8 +674,8 @@ public class DataSourcesResourceTest
);
// test dropped
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(TestDataSource.WIKI))
- .andReturn(ImmutableList.of(loadRule, dropRule))
+ EasyMock.expect(databaseRuleManager.getRulesSnapshot())
+ .andReturn(clusterDefaultRules(loadRule, dropRule))
.once();
EasyMock.replay(databaseRuleManager);
@@ -685,8 +687,8 @@ public class DataSourcesResourceTest
// test isn't dropped and no timeline found
EasyMock.reset(databaseRuleManager);
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(TestDataSource.WIKI))
- .andReturn(ImmutableList.of(loadRule, dropRule))
+ EasyMock.expect(databaseRuleManager.getRulesSnapshot())
+ .andReturn(clusterDefaultRules(loadRule, dropRule))
.once();
EasyMock.expect(inventoryView.getTimeline(new
TableDataSource(TestDataSource.WIKI)))
.andReturn(null)
@@ -717,8 +719,8 @@ public class DataSourcesResourceTest
}
};
EasyMock.reset(inventoryView, databaseRuleManager);
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(TestDataSource.WIKI))
- .andReturn(ImmutableList.of(loadRule, dropRule))
+ EasyMock.expect(databaseRuleManager.getRulesSnapshot())
+ .andReturn(clusterDefaultRules(loadRule, dropRule))
.once();
EasyMock.expect(inventoryView.getTimeline(new
TableDataSource(TestDataSource.WIKI)))
.andReturn(timeline)
@@ -755,8 +757,8 @@ public class DataSourcesResourceTest
null,
auditManager
);
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(TestDataSource.WIKI))
- .andReturn(ImmutableList.of(partialRule))
+ EasyMock.expect(databaseRuleManager.getRulesSnapshot())
+ .andReturn(clusterDefaultRules(partialRule))
.once();
EasyMock.expect(segmentsMetadataManager.getRecentDataSourcesSnapshot())
.andReturn(DataSourcesSnapshot.fromUsedSegments(ImmutableList.of()))
@@ -800,8 +802,8 @@ public class DataSourcesResourceTest
String interval = "2013-01-01T01:00:00Z/2013-01-01T02:00:00Z";
DataSegment segment = buildHandoffSegment(TestDataSource.WIKI,
Intervals.of(interval), "v1", 1);
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(TestDataSource.WIKI))
- .andReturn(ImmutableList.of(partialRule))
+ EasyMock.expect(databaseRuleManager.getRulesSnapshot())
+ .andReturn(clusterDefaultRules(partialRule))
.once();
EasyMock.expect(segmentsMetadataManager.getRecentDataSourcesSnapshot())
.andReturn(DataSourcesSnapshot.fromUsedSegments(ImmutableList.of()))
@@ -858,8 +860,8 @@ public class DataSourcesResourceTest
ImmutableList.of("other_daily")
);
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(TestDataSource.WIKI))
- .andReturn(ImmutableList.of(partialRule, dropRule))
+ EasyMock.expect(databaseRuleManager.getRulesSnapshot())
+ .andReturn(clusterDefaultRules(partialRule, dropRule))
.once();
EasyMock.expect(segmentsMetadataManager.getRecentDataSourcesSnapshot())
.andReturn(DataSourcesSnapshot.fromUsedSegments(ImmutableList.of(segment)))
@@ -907,8 +909,8 @@ public class DataSourcesResourceTest
ImmutableList.of("user_daily")
);
-
EasyMock.expect(databaseRuleManager.getRulesWithDefault(TestDataSource.WIKI))
- .andReturn(ImmutableList.of(partialRule, dropRule))
+ EasyMock.expect(databaseRuleManager.getRulesSnapshot())
+ .andReturn(clusterDefaultRules(partialRule, dropRule))
.once();
EasyMock.expect(segmentsMetadataManager.getRecentDataSourcesSnapshot())
.andReturn(DataSourcesSnapshot.fromUsedSegments(ImmutableList.of(segment)))
@@ -924,6 +926,17 @@ public class DataSourcesResourceTest
EasyMock.verify(inventoryView, databaseRuleManager,
segmentsMetadataManager);
}
+ /**
+ * Snapshot in which the given rules apply to every datasource.
+ */
+ private static RetentionRulesSnapshot clusterDefaultRules(Rule... rules)
+ {
+ return new RetentionRulesSnapshot(
+ Map.of(MetadataRuleManagerConfig.DEFAULT_RULE_NAME, List.of(rules)),
+ MetadataRuleManagerConfig.DEFAULT_RULE_NAME
+ );
+ }
+
private static DataSegment buildHandoffSegment(String dataSource, Interval
interval, String version, int partitionNumber)
{
return buildHandoffSegment(dataSource, interval, version, partitionNumber,
null);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]