This is an automated email from the ASF dual-hosted git repository.
kfaraz 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 4caa5ce4c44 minor: Add config and metrics to CostBasedAutoScaler
(#19697)
4caa5ce4c44 is described below
commit 4caa5ce4c44df2dadeab944b933c789da26f001c
Author: Kashif Faraz <[email protected]>
AuthorDate: Thu Jul 16 22:28:57 2026 +0530
minor: Add config and metrics to CostBasedAutoScaler (#19697)
Changes:
- Add config to skip scaling if cost drop percentage is less than a
threshold.
- This would help prevent the autoscaler from flapping.
- This is a temporary fix until we have stabilized the cost function.
- Add metrics for current cost and optimal cost in `CostBasedAutoScaler`
---
.../supervisor/autoscaler/CostBasedAutoScaler.java | 14 ++++++++++
.../autoscaler/CostBasedAutoScalerConfig.java | 30 +++++++++++++++++++---
.../autoscaler/CostBasedAutoScalerConfigTest.java | 5 +++-
3 files changed, 45 insertions(+), 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 b089df3bcd1..abada2f9b74 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
@@ -65,8 +65,10 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
public static final String LAG_WEIGHT_METRIC =
"task/autoScaler/costBased/lagWeight";
public static final String IDLE_WEIGHT_METRIC =
"task/autoScaler/costBased/idleWeight";
+ public static final String CURRENT_COST_METRIC =
"task/autoScaler/costBased/currentCost";
public static final String CURRENT_LAG_COST_METRIC =
"task/autoScaler/costBased/currentLagCost";
public static final String CURRENT_IDLE_COST_METRIC =
"task/autoScaler/costBased/currentIdleCost";
+ public static final String OPTIMAL_COST_METRIC =
"task/autoScaler/costBased/optimalCost";
public static final String OPTIMAL_LAG_COST_METRIC =
"task/autoScaler/costBased/optimalLagCost";
public static final String OPTIMAL_IDLE_COST_METRIC =
"task/autoScaler/costBased/optimalIdleCost";
public static final String OPTIMAL_TASK_COUNT_METRIC =
"task/autoScaler/costBased/optimalTaskCount";
@@ -307,6 +309,8 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
emitter.emit(getMetricBuilder().setMetric(IDLE_WEIGHT_METRIC,
config.getIdleWeight()));
emitter.emit(getMetricBuilder().setMetric(CURRENT_LAG_COST_METRIC,
currentCost.lagCost()));
emitter.emit(getMetricBuilder().setMetric(CURRENT_IDLE_COST_METRIC,
currentCost.idleCost()));
+ emitter.emit(getMetricBuilder().setMetric(CURRENT_COST_METRIC,
currentCost.totalCost()));
+ emitter.emit(getMetricBuilder().setMetric(OPTIMAL_COST_METRIC,
optimalCost.totalCost()));
// Emit avg rate and idle metrics only if they are available
if (metrics.getAvgProcessingRate() >= 0) {
@@ -326,6 +330,16 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
);
emitter.emit(getMetricBuilder().setMetric(OPTIMAL_LAG_COST_METRIC,
optimalCost.lagCost()));
emitter.emit(getMetricBuilder().setMetric(OPTIMAL_IDLE_COST_METRIC,
optimalCost.idleCost()));
+
+ final double costDropPercent
+ = 100.0 * (currentCost.totalCost() - optimalCost.totalCost()) /
currentCost.totalCost();
+ if (costDropPercent < config.getMinCostDropPercentForScaling()) {
+ log.info(
+ "Skipping scaling since cost drop percent[%.2f] is less than
required minCostDropPercentForScaling[%d]",
+ costDropPercent, config.getMinCostDropPercentForScaling()
+ );
+ return currentTaskCount;
+ }
}
// Scale-up is applied eagerly; scale-down may be deferred by
computeTaskCountForScaleAction().
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 785fa3257b5..b93c354b6a5 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
@@ -65,6 +65,7 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
private final Duration minScaleDownDelay;
private final boolean scaleDownDuringTaskRolloverOnly;
private final boolean usePollIdleRatio;
+ private final int minCostDropPercentForScaling;
/**
* Creates a new CostBasedAutoScalerConfig instance.
@@ -85,7 +86,8 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
@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("minCostDropPercentForScaling") Integer
minCostDropPercentForScaling
)
{
this.enableTaskAutoScaler = Configs.valueOrDefault(enableTaskAutoScaler,
false);
@@ -103,6 +105,7 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
this.minScaleDownDelay = Configs.valueOrDefault(minScaleDownDelay,
DEFAULT_MIN_SCALE_DOWN_DELAY);
this.scaleDownDuringTaskRolloverOnly =
Configs.valueOrDefault(scaleDownDuringTaskRolloverOnly, false);
this.usePollIdleRatio = Configs.valueOrDefault(usePollIdleRatio, true);
+ this.minCostDropPercentForScaling =
Configs.valueOrDefault(minCostDropPercentForScaling, 0);
if (this.enableTaskAutoScaler) {
Preconditions.checkNotNull(taskCountMax, "taskCountMax is required when
enableTaskAutoScaler is true");
@@ -285,6 +288,16 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
return usePollIdleRatio;
}
+ /**
+ * Minimum percentage drop from current cost that is required by the
auto-scaler
+ * to choose a new task count.
+ */
+ @JsonProperty
+ public int getMinCostDropPercentForScaling()
+ {
+ return minCostDropPercentForScaling;
+ }
+
@Override
public SupervisorTaskAutoScaler createAutoScaler(Supervisor supervisor,
SupervisorSpec spec, ServiceEmitter emitter)
{
@@ -316,6 +329,7 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
&& Objects.equals(minScaleDownDelay, that.minScaleDownDelay)
&& scaleDownDuringTaskRolloverOnly ==
that.scaleDownDuringTaskRolloverOnly
&& usePollIdleRatio == that.usePollIdleRatio
+ && minCostDropPercentForScaling == that.minCostDropPercentForScaling
&& Objects.equals(taskCountStart, that.taskCountStart)
&& Objects.equals(stopTaskCountRatio, that.stopTaskCountRatio);
}
@@ -338,7 +352,8 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
minScaleUpDelay,
minScaleDownDelay,
scaleDownDuringTaskRolloverOnly,
- usePollIdleRatio
+ usePollIdleRatio,
+ minCostDropPercentForScaling
);
}
@@ -361,6 +376,7 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
", minScaleDownDelay=" + minScaleDownDelay +
", scaleDownDuringTaskRolloverOnly=" +
scaleDownDuringTaskRolloverOnly +
", usePollIdleRatio=" + usePollIdleRatio +
+ ", minCostDropPercentForScaling=" + minCostDropPercentForScaling +
'}';
}
@@ -385,6 +401,7 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
private Duration minScaleDownDelay;
private Boolean scaleDownDuringTaskRolloverOnly;
private Boolean usePollIdleRatio;
+ private Integer minCostDropPercentForScaling;
private Builder()
{
@@ -480,6 +497,12 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
return this;
}
+ public Builder minCostDropPercentForScaling(int
minCostDropPercentForScaling)
+ {
+ this.minCostDropPercentForScaling = minCostDropPercentForScaling;
+ return this;
+ }
+
public CostBasedAutoScalerConfig build()
{
return new CostBasedAutoScalerConfig(
@@ -497,7 +520,8 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
minScaleUpDelay,
minScaleDownDelay,
scaleDownDuringTaskRolloverOnly,
- usePollIdleRatio
+ usePollIdleRatio,
+ 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 45771d782ab..295dd1efe88 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 class CostBasedAutoScalerConfigTest
+ " \"minScaleUpDelay\": \"PT5M\",\n"
+ " \"minScaleDownDelay\": \"PT10M\",\n"
+ " \"scaleDownDuringTaskRolloverOnly\": true,\n"
- + " \"usePollIdleRatio\": false\n"
+ + " \"usePollIdleRatio\": false,\n"
+ + " \"minCostDropPercentForScaling\": 10\n"
+ "}";
final CostBasedAutoScalerConfig config = mapper.readValue(json,
CostBasedAutoScalerConfig.class);
@@ -74,6 +75,7 @@ public class CostBasedAutoScalerConfigTest
Assert.assertFalse(config.isUsePollIdleRatio());
Assert.assertFalse(config.isUseTaskCountBoundariesOnScaleUp());
Assert.assertTrue(config.isUseTaskCountBoundariesOnScaleDown());
+ Assert.assertEquals(10, config.getMinCostDropPercentForScaling());
// Test serialization back to JSON
final String serialized = mapper.writeValueAsString(config);
@@ -112,6 +114,7 @@ public class CostBasedAutoScalerConfigTest
Assert.assertTrue(config.isUseTaskCountBoundariesOnScaleDown());
Assert.assertNull(config.getTaskCountStart());
Assert.assertNull(config.getStopTaskCountRatio());
+ Assert.assertEquals(0, config.getMinCostDropPercentForScaling());
}
@Test
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]