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]

Reply via email to