From f4029759ff21665af17b376b499b5f044b926542 Mon Sep 17 00:00:00 2001 From: Sasha Syrotenko Date: Fri, 3 Jul 2026 15:11:00 +0300 Subject: [PATCH 1/7] Introduce lag emergency --- .../autoscaler/CostBasedAutoScaler.java | 50 +++++++++++++- .../autoscaler/CostBasedAutoScalerConfig.java | 45 +++++++++++-- .../autoscaler/WeightedCostFunction.java | 29 ++++++++- .../CostBasedAutoScalerConfigTest.java | 18 ++++- .../autoscaler/CostBasedAutoScalerTest.java | 56 ++++++++++++++++ .../autoscaler/WeightedCostFunctionTest.java | 65 +++++++++++++++++++ 6 files changed, 255 insertions(+), 8 deletions(-) diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java index bf07545aaa53..bdf7c3391b66 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java @@ -141,7 +141,10 @@ private ServiceMetricEvent.Builder getMetricBuilder() public void start() { autoscalerExecutor.scheduleAtFixedRate( - supervisor.buildDynamicAllocationTask(this::computeTaskCountForScaleAction, () -> {}, emitter), + supervisor.buildDynamicAllocationTask( + this::computeTaskCountForScaleAction, () -> { + }, emitter + ), config.getScaleActionPeriodMillis(), config.getScaleActionPeriodMillis(), TimeUnit.MILLISECONDS @@ -207,6 +210,25 @@ public CostBasedAutoScalerConfig getConfig() return config; } + private boolean isCriticalLag(CostMetrics metrics) + { + final Long criticalLagThreshold = config.getCriticalLagThreshold(); + return metrics != null && criticalLagThreshold != null + && metrics.getAggregateLag() >= criticalLagThreshold * WeightedCostFunction.CRITICAL_LAG_TIER1_FRACTION; + } + + /** + * Whether the last collected metrics crossed {@link WeightedCostFunction#CRITICAL_LAG_TIER2_FRACTION} of + * {@link CostBasedAutoScalerConfig#getCriticalLagThreshold()}, meaning the argmin search should be + * skipped entirely in favor of jumping straight to the maximum task count. + */ + private boolean isEmergencyLag(CostMetrics metrics) + { + final Long criticalLagThreshold = config.getCriticalLagThreshold(); + return metrics != null && criticalLagThreshold != null + && metrics.getAggregateLag() >= criticalLagThreshold * WeightedCostFunction.CRITICAL_LAG_TIER2_FRACTION; + } + /** * Returns the lowest-cost task count given {@code metrics}, or {@link #CANNOT_COMPUTE} when * metrics are unusable. Returning the current task count means the current count is already @@ -244,6 +266,24 @@ int computeOptimalTaskCount(CostMetrics metrics) return currentTaskCount; } + final boolean criticalLag = isCriticalLag(metrics); + final boolean emergencyLag = isEmergencyLag(metrics); + if (emergencyLag) { + log.info( + "Supervisor[%s] aggregateLag[%.0f] crossed [%.0f%%] of criticalLagThreshold[%d]: skipping the argmin" + + " search and jumping straight to the maximum task count.", + supervisorId, metrics.getAggregateLag(), WeightedCostFunction.CRITICAL_LAG_TIER2_FRACTION * 100, + config.getCriticalLagThreshold() + ); + } else if (criticalLag) { + log.info( + "Supervisor[%s] aggregateLag[%.0f] crossed [%.0f%%] of criticalLagThreshold[%d]: widening scale-up" + + " candidates and maxing out the lag-amplification multiplier.", + supervisorId, metrics.getAggregateLag(), WeightedCostFunction.CRITICAL_LAG_TIER1_FRACTION * 100, + config.getCriticalLagThreshold() + ); + } + // Start with the current task count as optimal int optimalTaskCount = currentTaskCount; CostResult optimalCost = costFunction.computeCost(metrics, currentTaskCount, config); @@ -270,7 +310,7 @@ int computeOptimalTaskCount(CostMetrics metrics) int startIndex = 0; int endIndex = validTaskCounts.length - 1; - if (config.isUseTaskCountBoundariesOnScaleUp()) { + if (config.isUseTaskCountBoundariesOnScaleUp() && !criticalLag) { int currentTaskCountIndex = Arrays.binarySearch(validTaskCounts, currentTaskCount); endIndex = currentTaskCountIndex >= 0 ? Math.min(currentTaskCountIndex + BOUNDARY_LIMIT_IN_PARTITIONS_PER_TASK, endIndex) @@ -284,6 +324,12 @@ int computeOptimalTaskCount(CostMetrics metrics) : startIndex; } + // Emergency (tier 2) lag skips the argmin search entirely: evaluate only the maximum valid task count. + if (emergencyLag) { + startIndex = validTaskCounts.length - 1; + endIndex = validTaskCounts.length - 1; + } + for (int i = startIndex; i <= endIndex; ++i) { final int taskCount = validTaskCounts[i]; CostResult costResult = costFunction.computeCost(metrics, taskCount, config); diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java index fa4d547eae49..0d84d45a19f2 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java @@ -69,6 +69,7 @@ public class CostBasedAutoScalerConfig implements AutoScalerConfig private final Duration minScaleDownDelay; private final boolean scaleDownDuringTaskRolloverOnly; private final boolean usePollIdleRatio; + private final Long criticalLagThreshold; /** * Creates a new CostBasedAutoScalerConfig instance. @@ -93,7 +94,8 @@ public CostBasedAutoScalerConfig( @Nullable @JsonProperty("minScaleUpDelay") Duration minScaleUpDelay, @Nullable @JsonProperty("minScaleDownDelay") Duration minScaleDownDelay, @Nullable @JsonProperty("scaleDownDuringTaskRolloverOnly") Boolean scaleDownDuringTaskRolloverOnly, - @Nullable @JsonProperty("usePollIdleRatio") Boolean usePollIdleRatio + @Nullable @JsonProperty("usePollIdleRatio") Boolean usePollIdleRatio, + @Nullable @JsonProperty("criticalLagThreshold") Long criticalLagThreshold ) { this.enableTaskAutoScaler = enableTaskAutoScaler != null ? enableTaskAutoScaler : false; @@ -123,6 +125,12 @@ public CostBasedAutoScalerConfig( this.minScaleDownDelay = Configs.valueOrDefault(minScaleDownDelay, DEFAULT_MIN_SCALE_DELAY); this.scaleDownDuringTaskRolloverOnly = Configs.valueOrDefault(scaleDownDuringTaskRolloverOnly, false); this.usePollIdleRatio = Configs.valueOrDefault(usePollIdleRatio, true); + this.criticalLagThreshold = criticalLagThreshold; + + Preconditions.checkArgument( + criticalLagThreshold == null || criticalLagThreshold > 0, + "criticalLagThreshold must be > 0" + ); if (this.enableTaskAutoScaler) { Preconditions.checkNotNull(taskCountMax, "taskCountMax is required when enableTaskAutoScaler is true"); @@ -305,6 +313,24 @@ public boolean isUsePollIdleRatio() return usePollIdleRatio; } + /** + * Aggregate (sum-across-partitions) lag threshold driving a two-tier SLA-critical fast path, + * relative to {@link CostMetrics#getAggregateLag()}: + * + * {@code null} disables the feature. + */ + @JsonProperty + @Nullable + public Long getCriticalLagThreshold() + { + return criticalLagThreshold; + } + @Override public SupervisorTaskAutoScaler createAutoScaler(Supervisor supervisor, SupervisorSpec spec, ServiceEmitter emitter) { @@ -338,7 +364,8 @@ public boolean equals(Object o) && scaleDownDuringTaskRolloverOnly == that.scaleDownDuringTaskRolloverOnly && usePollIdleRatio == that.usePollIdleRatio && Objects.equals(taskCountStart, that.taskCountStart) - && Objects.equals(stopTaskCountRatio, that.stopTaskCountRatio); + && Objects.equals(stopTaskCountRatio, that.stopTaskCountRatio) + && Objects.equals(criticalLagThreshold, that.criticalLagThreshold); } @Override @@ -360,7 +387,8 @@ public int hashCode() minScaleUpDelay, minScaleDownDelay, scaleDownDuringTaskRolloverOnly, - usePollIdleRatio + usePollIdleRatio, + criticalLagThreshold ); } @@ -384,6 +412,7 @@ public String toString() ", minScaleDownDelay=" + minScaleDownDelay + ", scaleDownDuringTaskRolloverOnly=" + scaleDownDuringTaskRolloverOnly + ", usePollIdleRatio=" + usePollIdleRatio + + ", criticalLagThreshold=" + criticalLagThreshold + '}'; } @@ -409,6 +438,7 @@ public static class Builder private Duration minScaleDownDelay; private Boolean scaleDownDuringTaskRolloverOnly; private Boolean usePollIdleRatio; + private Long criticalLagThreshold; private Builder() { @@ -498,6 +528,12 @@ public Builder usePollIdleRatio(boolean usePollIdleRatio) return this; } + public Builder criticalLagThreshold(Long criticalLagThreshold) + { + this.criticalLagThreshold = criticalLagThreshold; + return this; + } + public Builder useTaskCountBoundariesOnScaleUp(boolean useTaskCountBoundariesOnScaleUp) { this.useTaskCountBoundariesOnScaleUp = useTaskCountBoundariesOnScaleUp; @@ -528,7 +564,8 @@ public CostBasedAutoScalerConfig build() minScaleUpDelay, minScaleDownDelay, scaleDownDuringTaskRolloverOnly, - usePollIdleRatio + usePollIdleRatio, + criticalLagThreshold ); } } diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunction.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunction.java index beaf0a5b9d55..c0d330edb526 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunction.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunction.java @@ -42,6 +42,26 @@ public class WeightedCostFunction */ static final double LAG_AMPLIFICATION_MULTIPLIER = 0.3; + /** + * Amplification multiplier used once aggregate lag crosses {@link #CRITICAL_LAG_TIER1_FRACTION} of + * {@link CostBasedAutoScalerConfig#getCriticalLagThreshold()} (tier 1 of the critical-lag fast path). + */ + static final double CRITICAL_LAG_AMPLIFICATION_MULTIPLIER = 6.0; + + /** + * Fraction of {@link CostBasedAutoScalerConfig#getCriticalLagThreshold()} at which tier 1 of the + * critical-lag fast path engages: amplification maxes out at {@link #CRITICAL_LAG_AMPLIFICATION_MULTIPLIER} + * and the scale-up candidate boundary is bypassed. + */ + static final double CRITICAL_LAG_TIER1_FRACTION = 0.75; + + /** + * Fraction of {@link CostBasedAutoScalerConfig#getCriticalLagThreshold()} at which tier 2 of the + * critical-lag fast path engages: the cost-minimization search is skipped entirely and the task + * count jumps straight to the maximum. + */ + static final double CRITICAL_LAG_TIER2_FRACTION = 0.95; + /** * Exponent (< 1) for sublinear busy redistribution in the idle projection: * busy grows as {@code (currentTaskCount / proposedTaskCount)^EXPONENT}, not linearly. @@ -108,12 +128,19 @@ public CostResult computeCost( // Lag recovery time is decreasing by adding tasks and increasing by ejecting tasks. // In case of increasing lag, we apply an amplification factor to reflect the urgency of addressing lag. // Caution: we rely only on the metrics, the real issues may be absolutely different, up to hardware failure. + // Once aggregate lag crosses CRITICAL_LAG_TIER1_FRACTION of criticalLagThreshold, the multiplier is + // maxed out at CRITICAL_LAG_AMPLIFICATION_MULTIPLIER (vs the default 0.3). final double lagRecoveryTime; if (metrics.getAggregateLag() <= 0) { lagRecoveryTime = 0; } else { final double lagPerPartition = metrics.getAggregateLag() / metrics.getPartitionCount(); - final double amplification = Math.max(1.0, 1.0 + LAG_AMPLIFICATION_MULTIPLIER * Math.log(lagPerPartition)); + final Long criticalLagThreshold = config.getCriticalLagThreshold(); + final boolean criticalLag = criticalLagThreshold != null + && metrics.getAggregateLag() >= criticalLagThreshold * CRITICAL_LAG_TIER1_FRACTION; + + final double amplificationMultiplier = criticalLag ? CRITICAL_LAG_AMPLIFICATION_MULTIPLIER : LAG_AMPLIFICATION_MULTIPLIER; + final double amplification = Math.max(1.0, 1.0 + amplificationMultiplier * Math.log(lagPerPartition)); final double adjustedProcessingRate = Math.max(avgProcessingRate, MIN_PROCESSING_RATE); lagRecoveryTime = metrics.getAggregateLag() * amplification / (proposedTaskCount * adjustedProcessingRate); } diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java index 6d09221b410a..c802840cf4df 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java @@ -54,7 +54,8 @@ public void testSerdeWithAllProperties() throws Exception + " \"minScaleUpDelay\": \"PT5M\",\n" + " \"minScaleDownDelay\": \"PT10M\",\n" + " \"scaleDownDuringTaskRolloverOnly\": true,\n" - + " \"usePollIdleRatio\": false\n" + + " \"usePollIdleRatio\": false,\n" + + " \"criticalLagThreshold\": 500000\n" + "}"; final CostBasedAutoScalerConfig config = mapper.readValue(json, CostBasedAutoScalerConfig.class); @@ -74,6 +75,7 @@ public void testSerdeWithAllProperties() throws Exception Assert.assertFalse(config.isUsePollIdleRatio()); Assert.assertFalse(config.isUseTaskCountBoundariesOnScaleUp()); Assert.assertTrue(config.isUseTaskCountBoundariesOnScaleDown()); + Assert.assertEquals(Long.valueOf(500000), config.getCriticalLagThreshold()); // Test serialization back to JSON final String serialized = mapper.writeValueAsString(config); @@ -112,6 +114,7 @@ public void testSerdeWithDefaults() throws Exception Assert.assertTrue(config.isUseTaskCountBoundariesOnScaleDown()); Assert.assertNull(config.getTaskCountStart()); Assert.assertNull(config.getStopTaskCountRatio()); + Assert.assertNull(config.getCriticalLagThreshold()); } @Test @@ -221,6 +224,7 @@ public void testBuilder() .minScaleDownDelay(Duration.standardMinutes(10)) .scaleDownDuringTaskRolloverOnly(true) .usePollIdleRatio(false) + .criticalLagThreshold(500000L) .build(); Assert.assertTrue(config.getEnableTaskAutoScaler()); @@ -238,6 +242,18 @@ public void testBuilder() Assert.assertEquals(Duration.standardMinutes(10), config.getMinScaleDownDelay()); Assert.assertTrue(config.isScaleDownOnTaskRolloverOnly()); Assert.assertFalse(config.isUsePollIdleRatio()); + Assert.assertEquals(Long.valueOf(500000), config.getCriticalLagThreshold()); + } + + @Test(expected = IllegalArgumentException.class) + public void testValidation_ZeroCriticalLagThreshold() + { + CostBasedAutoScalerConfig.builder() + .taskCountMax(100) + .taskCountMin(5) + .criticalLagThreshold(0L) + .enableTaskAutoScaler(true) + .build(); } @Test diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java index d773e20bfcb3..696160f2d90e 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java @@ -258,6 +258,62 @@ public void testComputeOptimalTaskCountLimitsTaskCountJumps() ); } + @Test + public void testCriticalLagThresholdBypassesScaleUpBoundary() + { + // aggregateLag = 100_000 * 100 = 10,000,000. With threshold=12,000,000: tier1=9,000,000 (crossed), + // tier2=11,400,000 (not crossed), so this exercises tier1 (boundary bypass) without triggering + // tier2's emergency jump-to-max. + final CostBasedAutoScalerConfig boundedScaleUpConfig = CostBasedAutoScalerConfig + .builder() + .taskCountMax(100) + .taskCountMin(1) + .enableTaskAutoScaler(true) + .lagWeight(1.0) + .idleWeight(0.0) + .useTaskCountBoundariesOnScaleUp(true) + .criticalLagThreshold(12_000_000L) + .build(); + final CostBasedAutoScaler scaler = createAutoScaler(boundedScaleUpConfig); + + Assert.assertEquals( + "Critical lag should bypass the scale-up boundary and jump straight to the argmin", + 100, + scaler.computeOptimalTaskCount(createMetrics(100_000.0, 10, 100, 0.25)) + ); + + // Below the threshold, the boundary still applies as usual. + Assert.assertEquals( + "Below criticalLagThreshold, the scale-up boundary still limits candidates", + 13, + scaler.computeOptimalTaskCount(createMetrics(10.0, 10, 100, 0.25)) + ); + } + + @Test + public void testEmergencyLagJumpsStraightToMaxTaskCount() + { + // aggregateLag = 100_000 * 500 = 50,000,000. With threshold=10,000,000: tier2=9,500,000 is + // comfortably crossed, so the argmin search is skipped entirely in favor of the maximum task count. + final CostBasedAutoScalerConfig config = CostBasedAutoScalerConfig + .builder() + .taskCountMax(500) + .taskCountMin(1) + .enableTaskAutoScaler(true) + .lagWeight(0.1) + .idleWeight(0.9) + .criticalLagThreshold(10_000_000L) + .build(); + final CostBasedAutoScaler scaler = createAutoScaler(config); + + // Idle-heavy weights would normally argue for scaling down, but emergency lag overrides that entirely. + Assert.assertEquals( + "Emergency lag should jump straight to the maximum task count regardless of idle-favoring weights", + 500, + scaler.computeOptimalTaskCount(createMetrics(100_000.0, 10, 500, 0.9)) + ); + } + @Test public void testExtractPollIdleRatio() { diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java index 6802692ade49..da4c68dbbf4b 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java @@ -357,6 +357,71 @@ public void testLagAmplificationAppliedUnconditionally() Assert.assertEquals("Lag amplification should increase lag recovery time", expected, costWithAmp, 0.0001); } + @Test + public void testCriticalLagThresholdMaxesOutAmplificationMultiplier() + { + int currentTaskCount = 10; + int proposedTaskCount = 10; + int partitionCount = 10; + double avgPartitionLag = 150.0; + double aggregateLag = avgPartitionLag * partitionCount; + + CostMetrics metrics = createMetrics(avgPartitionLag, currentTaskCount, partitionCount, 0.1); + + // aggregateLag sits at exactly tier1Fraction (75%) of this threshold. + long tier1Threshold = (long) (aggregateLag / WeightedCostFunction.CRITICAL_LAG_TIER1_FRACTION); + + CostBasedAutoScalerConfig noThreshold = CostBasedAutoScalerConfig.builder() + .taskCountMax(100) + .taskCountMin(1) + .enableTaskAutoScaler(true) + .lagWeight(1.0) + .idleWeight(0.0) + .build(); + CostBasedAutoScalerConfig belowTier1 = CostBasedAutoScalerConfig.builder() + .taskCountMax(100) + .taskCountMin(1) + .enableTaskAutoScaler(true) + .lagWeight(1.0) + .idleWeight(0.0) + .criticalLagThreshold(tier1Threshold + 100) + .build(); + CostBasedAutoScalerConfig atTier1 = CostBasedAutoScalerConfig.builder() + .taskCountMax(100) + .taskCountMin(1) + .enableTaskAutoScaler(true) + .lagWeight(1.0) + .idleWeight(0.0) + .criticalLagThreshold(tier1Threshold) + .build(); + + double costBelowTier1 = costFunction.computeCost(metrics, proposedTaskCount, belowTier1).totalCost(); + Assert.assertEquals( + "Below tier1, amplification uses the default multiplier", + costFunction.computeCost(metrics, proposedTaskCount, noThreshold).totalCost(), + costBelowTier1, + 0.0001 + ); + + double lagPerPartition = aggregateLag / partitionCount; + double criticalAmplification = + 1.0 + WeightedCostFunction.CRITICAL_LAG_AMPLIFICATION_MULTIPLIER * Math.log(lagPerPartition); + double expectedCriticalCost = + aggregateLag * criticalAmplification / (proposedTaskCount * WeightedCostFunction.MIN_PROCESSING_RATE); + + double costAtTier1 = costFunction.computeCost(metrics, proposedTaskCount, atTier1).totalCost(); + Assert.assertEquals( + "At/above tier1, the amplification multiplier maxes out at CRITICAL_LAG_AMPLIFICATION_MULTIPLIER", + expectedCriticalCost, + costAtTier1, + 0.0001 + ); + Assert.assertTrue( + "Critical-lag cost should exceed the default-multiplier cost for the same lag", + costAtTier1 > costBelowTier1 + ); + } + @Test public void testAmplificationGrowsWithLag() { From d555233117f628d4f3267870de6b79795bf8b022 Mon Sep 17 00:00:00 2001 From: Sasha Syrotenko Date: Fri, 3 Jul 2026 16:20:52 +0300 Subject: [PATCH 2/7] Checkstyle --- .../supervisor/autoscaler/CostBasedAutoScaler.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java index bdf7c3391b66..631881db7a75 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java @@ -141,10 +141,7 @@ private ServiceMetricEvent.Builder getMetricBuilder() public void start() { autoscalerExecutor.scheduleAtFixedRate( - supervisor.buildDynamicAllocationTask( - this::computeTaskCountForScaleAction, () -> { - }, emitter - ), + supervisor.buildDynamicAllocationTask(this::computeTaskCountForScaleAction, () -> {}, emitter), config.getScaleActionPeriodMillis(), config.getScaleActionPeriodMillis(), TimeUnit.MILLISECONDS From ffa2092f65c5d847869ff8318ef1cc3c37a2f6bf Mon Sep 17 00:00:00 2001 From: Sasha Syrotenko Date: Sat, 11 Jul 2026 09:03:59 +0300 Subject: [PATCH 3/7] Leave lag multiplier as lag budget orchestrating tool --- .../autoscaler/CostBasedAutoScalerConfig.java | 5 ++-- .../autoscaler/WeightedCostFunction.java | 11 ++++----- .../autoscaler/WeightedCostFunctionTest.java | 24 +++++++++---------- 3 files changed, 20 insertions(+), 20 deletions(-) diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java index 0d84d45a19f2..2e725923fc79 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java @@ -317,8 +317,9 @@ public boolean isUsePollIdleRatio() * Aggregate (sum-across-partitions) lag threshold driving a two-tier SLA-critical fast path, * relative to {@link CostMetrics#getAggregateLag()}: *
    - *
  • At 75% of this value, the lag-amplification multiplier maxes out at 6.0 (instead of the - * default 0.3), and the scale-up candidate search bypasses {@link #isUseTaskCountBoundariesOnScaleUp()}.
  • + *
  • At 75% of this value, the lag-amplification multiplier maxes out at 6.0 (instead of + * unamplified normal recovery), and the scale-up candidate search bypasses + * {@link #isUseTaskCountBoundariesOnScaleUp()}.
  • *
  • At 95% of this value, cost minimization is skipped entirely and the task count jumps * straight to the maximum.
  • *
diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunction.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunction.java index c0d330edb526..99484d971581 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunction.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunction.java @@ -37,10 +37,10 @@ public class WeightedCostFunction private static final Logger log = new Logger(WeightedCostFunction.class); /** - * Multiplier for a lag amplification factor; it was carefully chosen - * during extensive testing as the most balanced multiplier for high-lag recovery. + * Normal-path lag amplification multiplier. Critical-lag tiers provide the + * urgency amplification, so normal lag uses unamplified recovery time. */ - static final double LAG_AMPLIFICATION_MULTIPLIER = 0.3; + static final double LAG_AMPLIFICATION_MULTIPLIER = 0.0; /** * Amplification multiplier used once aggregate lag crosses {@link #CRITICAL_LAG_TIER1_FRACTION} of @@ -126,10 +126,9 @@ public CostResult computeCost( } // Lag recovery time is decreasing by adding tasks and increasing by ejecting tasks. - // In case of increasing lag, we apply an amplification factor to reflect the urgency of addressing lag. - // Caution: we rely only on the metrics, the real issues may be absolutely different, up to hardware failure. + // Critical lag uses extra amplification; normal lag uses raw recovery time so capacity cost remains meaningful. // Once aggregate lag crosses CRITICAL_LAG_TIER1_FRACTION of criticalLagThreshold, the multiplier is - // maxed out at CRITICAL_LAG_AMPLIFICATION_MULTIPLIER (vs the default 0.3). + // maxed out at CRITICAL_LAG_AMPLIFICATION_MULTIPLIER (vs unamplified normal recovery). final double lagRecoveryTime; if (metrics.getAggregateLag() <= 0) { lagRecoveryTime = 0; diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java index da4c68dbbf4b..888549f9219d 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java @@ -329,7 +329,7 @@ public void testIdleRatioWithMissingData() } @Test - public void testLagAmplificationAppliedUnconditionally() + public void testNormalLagCostUsesUnamplifiedRecoveryTime() { CostBasedAutoScalerConfig lagOnly = CostBasedAutoScalerConfig.builder() .taskCountMax(100) @@ -344,17 +344,15 @@ public void testLagAmplificationAppliedUnconditionally() int partitionCount = 10; double pollIdleRatio = 0.1; - // lagPerPartition = 150 * 10 / 10 = 150, amplification = 1 + 0.2 * ln(150) + // Normal lag uses raw recovery time; critical lag is tested separately below. CostMetrics metrics = createMetrics(150.0, currentTaskCount, partitionCount, pollIdleRatio); double costWithAmp = costFunction.computeCost(metrics, proposedTaskCount, lagOnly).totalCost(); double aggregateLag = 150.0 * partitionCount; - double lagPerPartition = aggregateLag / partitionCount; - double amplification = 1.0 + WeightedCostFunction.LAG_AMPLIFICATION_MULTIPLIER * Math.log(lagPerPartition); - double expected = aggregateLag * amplification / (proposedTaskCount * WeightedCostFunction.MIN_PROCESSING_RATE); + double expected = aggregateLag / (proposedTaskCount * WeightedCostFunction.MIN_PROCESSING_RATE); - Assert.assertEquals("Lag amplification should increase lag recovery time", expected, costWithAmp, 0.0001); + Assert.assertEquals("Normal lag cost should use raw recovery time", expected, costWithAmp, 0.0001); } @Test @@ -423,9 +421,9 @@ public void testCriticalLagThresholdMaxesOutAmplificationMultiplier() } @Test - public void testAmplificationGrowsWithLag() + public void testNormalLagCostScalesLinearlyWithLag() { - // Verify that higher lag produces proportionally higher cost due to log amplification + // Without normal-path amplification, cost grows linearly with lag. CostBasedAutoScalerConfig lagOnly = CostBasedAutoScalerConfig.builder() .taskCountMax(100) .taskCountMin(1) @@ -447,12 +445,14 @@ public void testAmplificationGrowsWithLag() Assert.assertTrue("Higher lag should produce higher cost", highCost > lowCost); - // The ratio of costs should be more than the ratio of raw lags (due to amplification) + // The ratio of costs matches the ratio of raw lags. double lagRatio = 10_000.0 / 100.0; double costRatio = highCost / lowCost; - Assert.assertTrue( - "Amplification should make cost grow faster than linear with lag", - costRatio > lagRatio + Assert.assertEquals( + "Normal lag cost should grow linearly with lag", + lagRatio, + costRatio, + 0.0001 ); } From dca6cf75f4d3a324ee530d5f30640a8947c9d85c Mon Sep 17 00:00:00 2001 From: Sasha Syrotenko Date: Mon, 13 Jul 2026 18:14:00 +0300 Subject: [PATCH 4/7] Use exact total lag for tier checks --- .../autoscaler/CostBasedAutoScaler.java | 4 ++ .../supervisor/autoscaler/CostMetrics.java | 7 +- .../CostBasedAutoScalerMockTest.java | 1 + .../autoscaler/CostBasedAutoScalerTest.java | 72 +++++++++++++++++++ .../autoscaler/WeightedCostFunctionTest.java | 4 +- 5 files changed, 85 insertions(+), 3 deletions(-) diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java index 631881db7a75..988f8d9fa08d 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScaler.java @@ -480,11 +480,14 @@ CostMetrics collectMetrics() final LagStats lagStats = supervisor.computeLagStats(); final double avgPartitionLag; + final double aggregateLag; if (lagStats == null) { log.debug("Lag stats unavailable for supervisorId [%s], skipping collection", supervisorId); avgPartitionLag = -1; + aggregateLag = -1; } else { avgPartitionLag = lagStats.getAvgLag(); + aggregateLag = lagStats.getTotalLag(); } final int currentTaskCount = supervisor.getIoConfig().getTaskCount(); @@ -500,6 +503,7 @@ CostMetrics collectMetrics() return new CostMetrics( avgPartitionLag, + aggregateLag, currentTaskCount, partitionCount, pollIdleRatio, diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostMetrics.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostMetrics.java index 33bf7eff33d5..f3e898497664 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostMetrics.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostMetrics.java @@ -41,6 +41,7 @@ public class CostMetrics public CostMetrics( double avgPartitionLag, + double aggregateLag, int currentTaskCount, int partitionCount, double pollIdleRatio, @@ -55,7 +56,7 @@ public CostMetrics( this.pollIdleRatio = pollIdleRatio; this.taskDurationSeconds = taskDurationSeconds; this.avgProcessingRate = avgProcessingRate; - this.aggregateLag = avgPartitionLag * partitionCount; + this.aggregateLag = aggregateLag; this.maxObservedRate = maxObservedRate; } @@ -89,7 +90,6 @@ public double getPollIdleRatio() /** * Returns the aggregated lag across all partitions. - * Pre-computed as avgPartitionLag * partitionCount. */ public double getAggregateLag() { @@ -142,6 +142,7 @@ public boolean equals(Object o) } CostMetrics that = (CostMetrics) o; return Double.compare(that.avgPartitionLag, avgPartitionLag) == 0 + && Double.compare(that.aggregateLag, aggregateLag) == 0 && currentTaskCount == that.currentTaskCount && partitionCount == that.partitionCount && Double.compare(that.pollIdleRatio, pollIdleRatio) == 0 @@ -155,6 +156,7 @@ public int hashCode() { return Objects.hash( avgPartitionLag, + aggregateLag, currentTaskCount, partitionCount, pollIdleRatio, @@ -169,6 +171,7 @@ public String toString() { return "CostMetrics{" + "avgPartitionLag=" + avgPartitionLag + + ", aggregateLag=" + aggregateLag + ", currentTaskCount=" + currentTaskCount + ", partitionCount=" + partitionCount + ", pollIdleRatio=" + pollIdleRatio + diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerMockTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerMockTest.java index fa86fbf1d3cb..3430842517dd 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerMockTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerMockTest.java @@ -367,6 +367,7 @@ private void setupMocksForMetricsCollection( { CostMetrics metrics = new CostMetrics( avgLag, + avgLag * PARTITION_COUNT, taskCount, PARTITION_COUNT, pollIdleRatio, diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java index 696160f2d90e..104857241ca9 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java @@ -290,6 +290,28 @@ public void testCriticalLagThresholdBypassesScaleUpBoundary() ); } + @Test + public void testCriticalLagThresholdUsesExactAggregateLag() + { + final CostBasedAutoScalerConfig boundedScaleUpConfig = CostBasedAutoScalerConfig + .builder() + .taskCountMax(1_000) + .taskCountMin(1) + .enableTaskAutoScaler(true) + .lagWeight(1.0) + .idleWeight(0.0) + .useTaskCountBoundariesOnScaleUp(true) + .criticalLagThreshold(1_000L) + .build(); + final CostBasedAutoScaler scaler = createAutoScaler(boundedScaleUpConfig); + + Assert.assertEquals( + "Exact aggregate lag should engage tier 1 even when integer average lag is zero", + 1_000, + scaler.computeOptimalTaskCount(createMetrics(0.0, 999.0, 10, 1_000, 0.25)) + ); + } + @Test public void testEmergencyLagJumpsStraightToMaxTaskCount() { @@ -627,6 +649,35 @@ public void testScalingActionSkippedWhenMovingAverageRateUnavailable() ); } + @Test + public void testCollectMetricsPreservesExactAggregateLag() + { + final SupervisorSpec spec = Mockito.mock(SupervisorSpec.class); + final SeekableStreamSupervisor supervisor = Mockito.mock(SeekableStreamSupervisor.class); + final ServiceEmitter emitter = Mockito.mock(ServiceEmitter.class); + final SeekableStreamSupervisorIOConfig ioConfig = Mockito.mock(SeekableStreamSupervisorIOConfig.class); + + when(spec.getId()).thenReturn("test-supervisor"); + when(spec.getDataSources()).thenReturn(List.of("test-datasource")); + when(spec.isSuspended()).thenReturn(false); + when(supervisor.getIoConfig()).thenReturn(ioConfig); + when(ioConfig.getStream()).thenReturn("test-stream"); + when(ioConfig.getTaskDuration()).thenReturn(Duration.standardHours(1)); + when(ioConfig.getTaskCount()).thenReturn(10); + when(supervisor.getPartitionCount()).thenReturn(1_000); + when(supervisor.computeLagStats()).thenReturn(new LagStats(999, 999, 0)); + when(supervisor.getStats()).thenReturn(Collections.emptyMap()); + + final CostBasedAutoScalerConfig config = CostBasedAutoScalerConfig.builder() + .taskCountMax(1_000) + .taskCountMin(1) + .enableTaskAutoScaler(true) + .build(); + final CostBasedAutoScaler scaler = new CostBasedAutoScaler(supervisor, config, spec, emitter); + + Assert.assertEquals(999.0, scaler.collectMetrics().getAggregateLag(), 0.0); + } + @Test public void testCollectMetricsTracksMaxProcessingRateOnlyWhenPollIdleRatioDisabled() { @@ -709,6 +760,27 @@ private CostMetrics createMetrics( { return new CostMetrics( avgPartitionLag, + avgPartitionLag * partitionCount, + currentTaskCount, + partitionCount, + pollIdleRatio, + 3600, + 1000.0, + 0. + ); + } + + private CostMetrics createMetrics( + double avgPartitionLag, + double aggregateLag, + int currentTaskCount, + int partitionCount, + double pollIdleRatio + ) + { + return new CostMetrics( + avgPartitionLag, + aggregateLag, currentTaskCount, partitionCount, pollIdleRatio, diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java index 888549f9219d..c3aaf8e01ae1 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/WeightedCostFunctionTest.java @@ -606,6 +606,7 @@ private CostMetrics createMetrics( { return new CostMetrics( avgPartitionLag, + avgPartitionLag * partitionCount, currentTaskCount, partitionCount, pollIdleRatio, @@ -621,7 +622,7 @@ private CostMetrics createMetricsWithMaxObservedRate( double pollIdleRatio ) { - return new CostMetrics(0.0, 10, 100, pollIdleRatio, 3600, avgProcessingRate, maxObservedRate); + return new CostMetrics(0.0, 0.0, 10, 100, pollIdleRatio, 3600, avgProcessingRate, maxObservedRate); } private CostMetrics createMetricsWithRate( @@ -634,6 +635,7 @@ private CostMetrics createMetricsWithRate( { return new CostMetrics( avgPartitionLag, + avgPartitionLag * partitionCount, currentTaskCount, partitionCount, pollIdleRatio, From ad6e74a6ac768e0ff25452bfdce326f01ee4b73a Mon Sep 17 00:00:00 2001 From: Sasha Syrotenko Date: Wed, 15 Jul 2026 12:21:10 +0300 Subject: [PATCH 5/7] Enhance default weights --- docs/ingestion/supervisor.md | 6 ++--- .../autoscaler/CostBasedAutoScalerConfig.java | 13 +++++------ .../CostBasedAutoScalerConfigTest.java | 23 ++++++++++--------- .../autoscaler/CostBasedAutoScalerTest.java | 2 +- 4 files changed, 22 insertions(+), 22 deletions(-) diff --git a/docs/ingestion/supervisor.md b/docs/ingestion/supervisor.md index b642aefec45e..582893f3713f 100644 --- a/docs/ingestion/supervisor.md +++ b/docs/ingestion/supervisor.md @@ -208,13 +208,13 @@ The following table outlines the configuration properties related to the `costBa | Property | Description | Required | Default | |----------|-------------|----------|---------------------------| -|`scaleActionPeriodMillis`|How often, in milliseconds, Druid evaluates whether to scale.|No| `600000` (10 min) | +|`scaleActionPeriodMillis`|How often, in milliseconds, Druid evaluates whether to scale.|No| `120000` (2 min) | |`lagWeight`|How much weight to give the lag cost relative to the idle cost. Higher values make the autoscaler more aggressive about adding tasks to drain backlog.|No| `0.4` | |`idleWeight`|How much weight to give the idle cost relative to the lag cost. Higher values make the autoscaler more aggressive about removing over-provisioned tasks.|No| `0.6` | |`useTaskCountBoundariesOnScaleUp`|Limits scale-up to a small step relative to the current task count, preventing large jumps. Disable to allow the autoscaler to jump directly to any task count.|No| `false` | |`useTaskCountBoundariesOnScaleDown`|Limits scale-down to a small step relative to the current task count, preventing large drops. Disable to allow the autoscaler to drop directly to any task count.|No| `true` | -|`minScaleUpDelay`|Minimum cooldown after a scale-up before the next scale-up is allowed. Specified as an ISO-8601 duration.|No| `scaleActionPeriodMillis` | -|`minScaleDownDelay`|Minimum cooldown after a scale-down before the next scale-down is allowed. Specified as an ISO-8601 duration.|No| `PT30M` | +|`minScaleUpDelay`|Minimum cooldown after a scale-up before the next scale-up is allowed. Specified as an ISO-8601 duration.|No| `PT15M` | +|`minScaleDownDelay`|Minimum cooldown after a scale-down before the next scale-down is allowed. Specified as an ISO-8601 duration.|No| `PT20M` | |`scaleDownDuringTaskRolloverOnly`|If `true`, scale-down actions are deferred until the next task rollover. This avoids disrupting in-progress ingestion.|No| `false` | The following example shows a supervisor spec with `costBased` autoscaler: diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java index 2e725923fc79..734cb6733c5e 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java @@ -47,10 +47,12 @@ public class CostBasedAutoScalerConfig implements AutoScalerConfig { private static final EmittingLogger LOG = new EmittingLogger(CostBasedAutoScalerConfig.class); - static final long DEFAULT_SCALE_ACTION_PERIOD_MILLIS = 10 * 60 * 1000; // 10 minutes + static final long DEFAULT_SCALE_ACTION_PERIOD_MILLIS = 2 * 60 * 1000; // 2 minutes + static final Duration DEFAULT_MIN_SCALE_UP_DELAY = Duration.millis(15 * 60 * 1000); // 15 minutes + static final Duration DEFAULT_MIN_SCALE_DOWN_DELAY = Duration.millis(20 * 60 * 1000); // 20 minutes + static final double DEFAULT_LAG_WEIGHT = 0.4; static final double DEFAULT_IDLE_WEIGHT = 0.6; - static final Duration DEFAULT_MIN_SCALE_DELAY = Duration.millis(DEFAULT_SCALE_ACTION_PERIOD_MILLIS * 3); private final boolean enableTaskAutoScaler; private final int taskCountMax; @@ -118,11 +120,8 @@ public CostBasedAutoScalerConfig( ); this.useTaskCountBoundariesOnScaleUp = Configs.valueOrDefault(useTaskCountBoundariesOnScaleUp, false); this.useTaskCountBoundariesOnScaleDown = Configs.valueOrDefault(useTaskCountBoundariesOnScaleDown, true); - this.minScaleUpDelay = Configs.valueOrDefault( - minScaleUpDelay, - Duration.millis(this.minTriggerScaleActionFrequencyMillis) - ); - this.minScaleDownDelay = Configs.valueOrDefault(minScaleDownDelay, DEFAULT_MIN_SCALE_DELAY); + this.minScaleUpDelay = Configs.valueOrDefault(minScaleUpDelay, DEFAULT_MIN_SCALE_UP_DELAY); + this.minScaleDownDelay = Configs.valueOrDefault(minScaleDownDelay, DEFAULT_MIN_SCALE_DOWN_DELAY); this.scaleDownDuringTaskRolloverOnly = Configs.valueOrDefault(scaleDownDuringTaskRolloverOnly, false); this.usePollIdleRatio = Configs.valueOrDefault(usePollIdleRatio, true); this.criticalLagThreshold = criticalLagThreshold; diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java index c802840cf4df..d04125caf3a9 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java @@ -27,7 +27,8 @@ import static org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostBasedAutoScalerConfig.DEFAULT_IDLE_WEIGHT; import static org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostBasedAutoScalerConfig.DEFAULT_LAG_WEIGHT; -import static org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostBasedAutoScalerConfig.DEFAULT_MIN_SCALE_DELAY; +import static org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostBasedAutoScalerConfig.DEFAULT_MIN_SCALE_DOWN_DELAY; +import static org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostBasedAutoScalerConfig.DEFAULT_MIN_SCALE_UP_DELAY; import static org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostBasedAutoScalerConfig.DEFAULT_SCALE_ACTION_PERIOD_MILLIS; import static org.apache.druid.indexing.seekablestream.supervisor.autoscaler.WeightedCostFunction.OPTIMAL_TASK_IDLE_RATIO; @@ -106,8 +107,8 @@ public void testSerdeWithDefaults() throws Exception Assert.assertEquals(DEFAULT_IDLE_WEIGHT, config.getIdleWeight(), 0.001); Assert.assertEquals(OPTIMAL_TASK_IDLE_RATIO, config.getOptimalTaskIdleRatio(), 0.001); // minScaleUpDelay and minScaleDownDelay each have their own independent default - Assert.assertEquals(Duration.millis(DEFAULT_SCALE_ACTION_PERIOD_MILLIS), config.getMinScaleUpDelay()); - Assert.assertEquals(DEFAULT_MIN_SCALE_DELAY, config.getMinScaleDownDelay()); + Assert.assertEquals(DEFAULT_MIN_SCALE_UP_DELAY, config.getMinScaleUpDelay()); + Assert.assertEquals(DEFAULT_MIN_SCALE_DOWN_DELAY, config.getMinScaleDownDelay()); Assert.assertFalse(config.isScaleDownOnTaskRolloverOnly()); Assert.assertTrue(config.isUsePollIdleRatio()); Assert.assertFalse(config.isUseTaskCountBoundariesOnScaleUp()); @@ -264,8 +265,8 @@ public void testScaleDelayDefaults() throws Exception .taskCountMax(10) .taskCountMin(1) .build(); - Assert.assertEquals(Duration.millis(DEFAULT_SCALE_ACTION_PERIOD_MILLIS), defaults.getMinScaleUpDelay()); - Assert.assertEquals(DEFAULT_MIN_SCALE_DELAY, defaults.getMinScaleDownDelay()); + Assert.assertEquals(DEFAULT_MIN_SCALE_UP_DELAY, defaults.getMinScaleUpDelay()); + Assert.assertEquals(DEFAULT_MIN_SCALE_DOWN_DELAY, defaults.getMinScaleDownDelay()); // Only minScaleUpDelay set: up uses explicit value, down uses its default CostBasedAutoScalerConfig upOnly = CostBasedAutoScalerConfig.builder() @@ -274,7 +275,7 @@ public void testScaleDelayDefaults() throws Exception .minScaleUpDelay(Duration.standardMinutes(5)) .build(); Assert.assertEquals(Duration.standardMinutes(5), upOnly.getMinScaleUpDelay()); - Assert.assertEquals(DEFAULT_MIN_SCALE_DELAY, upOnly.getMinScaleDownDelay()); + Assert.assertEquals(DEFAULT_MIN_SCALE_DOWN_DELAY, upOnly.getMinScaleDownDelay()); // Only minScaleDownDelay set: down uses explicit value, up uses its own default (does not fall back to down) CostBasedAutoScalerConfig downOnly = CostBasedAutoScalerConfig.builder() @@ -282,7 +283,7 @@ public void testScaleDelayDefaults() throws Exception .taskCountMin(1) .minScaleDownDelay(Duration.standardMinutes(20)) .build(); - Assert.assertEquals(Duration.millis(DEFAULT_SCALE_ACTION_PERIOD_MILLIS), downOnly.getMinScaleUpDelay()); + Assert.assertEquals(DEFAULT_MIN_SCALE_UP_DELAY, downOnly.getMinScaleUpDelay()); Assert.assertEquals(Duration.standardMinutes(20), downOnly.getMinScaleDownDelay()); // Both set: serde roundtrip preserves values @@ -306,8 +307,8 @@ public void testScaleDelayDefaults() throws Exception public void testMinTriggerScaleActionFrequencyMillisSerdeCompat() throws Exception { final long defaultMinTriggerMillis = DEFAULT_SCALE_ACTION_PERIOD_MILLIS; - final Duration defaultUp = Duration.millis(DEFAULT_SCALE_ACTION_PERIOD_MILLIS); - final Duration defaultDown = DEFAULT_MIN_SCALE_DELAY; + final Duration defaultUp = DEFAULT_MIN_SCALE_UP_DELAY; + final Duration defaultDown = DEFAULT_MIN_SCALE_DOWN_DELAY; // Backwards-compat: nothing set -> everything uses its own default. { @@ -330,7 +331,7 @@ public void testMinTriggerScaleActionFrequencyMillisSerdeCompat() throws Excepti CostBasedAutoScalerConfig.class ); Assert.assertEquals(900_000L, config.getMinTriggerScaleActionFrequencyMillis()); - Assert.assertEquals(Duration.millis(900_000L), config.getMinScaleUpDelay()); + Assert.assertEquals(defaultUp, config.getMinScaleUpDelay()); Assert.assertEquals(defaultDown, config.getMinScaleDownDelay()); assertRoundTrips(config); } @@ -387,7 +388,7 @@ public void testMinTriggerScaleActionFrequencyMillisSerdeCompat() throws Excepti CostBasedAutoScalerConfig.class ); Assert.assertEquals(900_000L, config.getMinTriggerScaleActionFrequencyMillis()); - Assert.assertEquals(Duration.millis(900_000L), config.getMinScaleUpDelay()); + Assert.assertEquals(defaultUp, config.getMinScaleUpDelay()); Assert.assertEquals(Duration.standardMinutes(15), config.getMinScaleDownDelay()); assertRoundTrips(config); } diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java index 104857241ca9..c0a8ac57b29f 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerTest.java @@ -590,7 +590,7 @@ public void testComputeTaskCountForRolloverAndConfigProperties() .enableTaskAutoScaler(true) .build(); Assert.assertEquals( - CostBasedAutoScalerConfig.DEFAULT_MIN_SCALE_DELAY, + CostBasedAutoScalerConfig.DEFAULT_MIN_SCALE_DOWN_DELAY, cfgWithDefaults.getMinScaleDownDelay() ); Assert.assertFalse(cfgWithDefaults.isScaleDownOnTaskRolloverOnly()); From e37f45f46e1e1c84fe1aee36127bcc0c15bf66c4 Mon Sep 17 00:00:00 2001 From: Sasha Syrotenko Date: Thu, 16 Jul 2026 13:31:55 +0300 Subject: [PATCH 6/7] Defaults --- .../supervisor/autoscaler/CostBasedAutoScalerConfig.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java index de995413d475..fac857ecb698 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java @@ -46,8 +46,8 @@ public class CostBasedAutoScalerConfig implements AutoScalerConfig { static final double DEFAULT_LAG_WEIGHT = 0.4; static final double DEFAULT_IDLE_WEIGHT = 0.6; - static final Duration DEFAULT_MIN_SCALE_UP_DELAY = Duration.standardMinutes(10); - static final Duration DEFAULT_MIN_SCALE_DOWN_DELAY = Duration.standardMinutes(30); + static final Duration DEFAULT_MIN_SCALE_UP_DELAY = Duration.standardMinutes(15); + static final Duration DEFAULT_MIN_SCALE_DOWN_DELAY = Duration.standardMinutes(20); static final Duration DEFAULT_SCALE_ACTION_PERIOD = Duration.standardMinutes(2); private final boolean enableTaskAutoScaler; From a9c1c1df8bf223e9b1d1b4273f60b947a55ad425 Mon Sep 17 00:00:00 2001 From: Sasha Syrotenko Date: Fri, 17 Jul 2026 11:22:44 +0300 Subject: [PATCH 7/7] Post-merge fix --- .../supervisor/autoscaler/CostBasedAutoScalerConfig.java | 4 ++-- .../supervisor/autoscaler/CostBasedAutoScalerConfigTest.java | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java index 5e622631aa96..47aadf644328 100644 --- a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java +++ b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfig.java @@ -89,7 +89,7 @@ public CostBasedAutoScalerConfig( @Nullable @JsonProperty("minScaleDownDelay") Duration minScaleDownDelay, @Nullable @JsonProperty("scaleDownDuringTaskRolloverOnly") Boolean scaleDownDuringTaskRolloverOnly, @Nullable @JsonProperty("usePollIdleRatio") Boolean usePollIdleRatio, - @Nullable @JsonProperty("criticalLagThreshold") Long criticalLagThreshold + @Nullable @JsonProperty("criticalLagThreshold") Long criticalLagThreshold, @Nullable @JsonProperty("minCostDropPercentForScaling") Integer minCostDropPercentForScaling ) { @@ -558,7 +558,7 @@ public CostBasedAutoScalerConfig build() minScaleDownDelay, scaleDownDuringTaskRolloverOnly, usePollIdleRatio, - criticalLagThreshold + criticalLagThreshold, minCostDropPercentForScaling ); } diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java index 9f0c3b303e25..472a55320920 100644 --- a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java +++ b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/autoscaler/CostBasedAutoScalerConfigTest.java @@ -56,7 +56,7 @@ public void testSerdeWithAllProperties() throws Exception + " \"minScaleDownDelay\": \"PT10M\",\n" + " \"scaleDownDuringTaskRolloverOnly\": true,\n" + " \"usePollIdleRatio\": false,\n" - + " \"criticalLagThreshold\": 500000\n" + + " \"criticalLagThreshold\": 500000,\n" + " \"minCostDropPercentForScaling\": 10\n" + "}";