This is an automated email from the ASF dual-hosted git repository. Fly-Style pushed a commit to branch cba-lag-emergency in repository https://gitbox.apache.org/repos/asf/druid.git
commit ad6e74a6ac768e0ff25452bfdce326f01ee4b73a Author: Sasha Syrotenko <[email protected]> AuthorDate: Wed Jul 15 12:21:10 2026 +0300 Enhance default weights --- docs/ingestion/supervisor.md | 6 +++--- .../autoscaler/CostBasedAutoScalerConfig.java | 13 ++++++------ .../autoscaler/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 b642aefec45..582893f3713 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 2e725923fc7..734cb6733c5 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 class CostBasedAutoScalerConfig implements AutoScalerConfig ); 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 c802840cf4d..d04125caf3a 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 org.junit.Test; 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 class CostBasedAutoScalerConfigTest 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 class CostBasedAutoScalerConfigTest .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 class CostBasedAutoScalerConfigTest .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 class CostBasedAutoScalerConfigTest .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 class CostBasedAutoScalerConfigTest 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 class CostBasedAutoScalerConfigTest 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 class CostBasedAutoScalerConfigTest 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 104857241ca..c0a8ac57b29 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 class CostBasedAutoScalerTest .enableTaskAutoScaler(true) .build(); Assert.assertEquals( - CostBasedAutoScalerConfig.DEFAULT_MIN_SCALE_DELAY, + CostBasedAutoScalerConfig.DEFAULT_MIN_SCALE_DOWN_DELAY, cfgWithDefaults.getMinScaleDownDelay() ); Assert.assertFalse(cfgWithDefaults.isScaleDownOnTaskRolloverOnly()); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
