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 f37486d51b6 feat: Add lag/idle plot simulator for cost-based
autoscaler (#19687)
f37486d51b6 is described below
commit f37486d51b6e8ac21ca3a536b38826b7d6edf5bb
Author: Kashif Faraz <[email protected]>
AuthorDate: Tue Sep 8 11:48:39 2026 +0530
feat: Add lag/idle plot simulator for cost-based autoscaler (#19687)
Changes currently in this PR:
- UI code is generated by Claude and may contain mistakes
- Add a simulate API currently supported for Kafka supervisors only
- This API creates a `CostBasedAutoScaler` in "simulate" mode and generates
the optimal task count for various input values of lag (based on criticalLag)
- Add a UI panel in the supervisor dialog which shows up only for "kafka"
supervisors
- Add a single plot between task count vs lag
---
.../embedded/indexing/KafkaClusterMetricsTest.java | 1 +
.../overlord/supervisor/SupervisorManager.java | 101 +++++++
.../overlord/supervisor/SupervisorResource.java | 37 ++-
.../supervisor/SeekableStreamSupervisor.java | 17 +-
.../supervisor/autoscaler/CostBasedAutoScaler.java | 140 +++++++---
.../autoscaler/CostBasedAutoScalerConfig.java | 28 +-
.../overlord/supervisor/SupervisorManagerTest.java | 75 +++++
.../auto-scaler-panel/auto-scaler-panel.scss | 68 +++++
.../auto-scaler-panel/auto-scaler-panel.spec.tsx | 41 +++
.../auto-scaler-panel/auto-scaler-panel.tsx | 304 +++++++++++++++++++++
.../supervisor-table-action-dialog.tsx | 19 +-
.../views/supervisors-view/supervisors-view.tsx | 11 +-
12 files changed, 794 insertions(+), 48 deletions(-)
diff --git
a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaClusterMetricsTest.java
b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaClusterMetricsTest.java
index 49be8522ca0..fc38808f41a 100644
---
a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaClusterMetricsTest.java
+++
b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaClusterMetricsTest.java
@@ -438,6 +438,7 @@ public class KafkaClusterMetricsTest extends
EmbeddedClusterTestBase
.withConsumerProperties(kafkaServer.consumerProperties())
.withTaskCount(taskCount)
)
+ .withContext(Map.of("useConcurrentLocks", true))
.withId(supervisorId)
.build(dataSource, TOPIC);
}
diff --git
a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManager.java
b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManager.java
index 507081c8340..3510ae7cfb6 100644
---
a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManager.java
+++
b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManager.java
@@ -26,9 +26,11 @@ import com.google.common.base.Preconditions;
import com.google.common.collect.ImmutableMap;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.inject.Inject;
+import org.apache.druid.common.config.Configs;
import org.apache.druid.common.guava.FutureUtils;
import org.apache.druid.common.utils.IdUtils;
import org.apache.druid.error.DruidException;
+import org.apache.druid.error.InternalServerError;
import org.apache.druid.error.InvalidInput;
import org.apache.druid.error.NotFound;
import org.apache.druid.guice.annotations.Json;
@@ -40,6 +42,9 @@ import
org.apache.druid.indexing.seekablestream.SeekableStreamDataSourceMetadata
import org.apache.druid.indexing.seekablestream.supervisor.BoundedStreamConfig;
import
org.apache.druid.indexing.seekablestream.supervisor.SeekableStreamSupervisor;
import
org.apache.druid.indexing.seekablestream.supervisor.SeekableStreamSupervisorSpec;
+import
org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostBasedAutoScaler;
+import
org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostBasedAutoScalerConfig;
+import
org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostMetrics;
import org.apache.druid.java.util.common.IAE;
import org.apache.druid.java.util.common.ISE;
import org.apache.druid.java.util.common.Pair;
@@ -59,6 +64,7 @@ import java.util.Collection;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.OptionalInt;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
@@ -640,6 +646,101 @@ public class SupervisorManager implements
SupervisorStatsProvider
}
}
+ /**
+ * Simulates the effects of the {@code costBased} auto-scaler by computing
the
+ * optimal task count under various values of aggregate lag.
+ *
+ * @return Map containing a single entry with key {@code "data"} and value as
+ * the simulation rows.
+ */
+ public Map<String, Object> simulateAutoscaling(
+ String supervisorId,
+ CostBasedAutoScalerConfig config,
+ int maxProcessingRatePerTask,
+ @Nullable Integer requestedTaskCount
+ )
+ {
+ if (!config.getEnableTaskAutoScaler()) {
+ throw InvalidInput.exception("Cannot simulate autoscaling since
'enableTaskAutoScaler' is false");
+ }
+
+ // Validate that this is a SeekableStreamSupervisor
+ final Pair<SeekableStreamSupervisor, SeekableStreamSupervisorSpec>
supervisorPair =
+ getSupervisorOfType(
+ supervisorId,
+ SeekableStreamSupervisor.class,
+ SeekableStreamSupervisorSpec.class,
+ "simulateAutoscaling"
+ );
+
+ // Validate the inputs
+ final long criticalLag =
Configs.valueOrDefault(config.getCriticalLagThreshold(), 1_000_000);
+ InvalidInput.conditionalException(
+ criticalLag >= 1000,
+ "Value of critical lag[%d] must be 1000 or more",
+ criticalLag
+ );
+ InvalidInput.conditionalException(
+ maxProcessingRatePerTask >= 100,
+ "Value of maxProcessingRatePerTask[%d] must be 100 events per second
or more",
+ maxProcessingRatePerTask
+ );
+ InvalidInput.conditionalException(
+ requestedTaskCount == null
+ || (requestedTaskCount >= config.getTaskCountMin() &&
requestedTaskCount <= config.getTaskCountMax()),
+ "Value of currentTaskCount[%d] must be within taskCountMin[%d] and
taskCountMax[%d]",
+ requestedTaskCount, config.getTaskCountMin(), config.getTaskCountMax()
+ );
+
+ // Simulate from the supervisor's live task count unless the caller pins
one.
+ final SeekableStreamSupervisorSpec supervisorSpec =
Objects.requireNonNull(supervisorPair.rhs);
+ final int currentTaskCount = supervisorSpec.getIoConfig().getTaskCount();
+ final int simulationTaskCount = Math.max(
+ config.getTaskCountMin(),
+ Math.min(
+ Configs.valueOrDefault(requestedTaskCount, currentTaskCount),
+ config.getTaskCountMax()
+ )
+ );
+
+ // Use the partition count and task duration from the supervisor spec
+ final int partitionCount =
Objects.requireNonNull(supervisorPair.lhs).getKnownPartitionCount();
+ if (partitionCount <= 0) {
+ throw InternalServerError.exception(
+ "Cannot simulate autoscaling since partition count for
supervisor[%s] is unknown."
+ + " Retry once the supervisor has discovered the current partition
count from stream.",
+ supervisorId
+ );
+ }
+ final long taskDurationSeconds =
supervisorSpec.getIoConfig().getTaskDuration().getStandardSeconds();
+
+ // Assume that the tasks are fully used since there is some lag
+ final double idleRatio = config.getOptimalTaskIdleRatio();
+
+ // Invoke the cost function for lag in the range [0, 2 *
criticalLagThreshold)
+ final Object[] rows = new Object[200];
+ final long lagStepSize = criticalLag / 100;
+ final CostBasedAutoScaler autoscaleSimulator =
CostBasedAutoScaler.createSimulator(config, supervisorId);
+ for (int i = 0; i < 200; ++i) {
+ final double observedAggregateLag = (double) lagStepSize * i;
+ final CostMetrics costMetrics = new CostMetrics(
+ observedAggregateLag / partitionCount,
+ observedAggregateLag,
+ simulationTaskCount,
+ partitionCount,
+ idleRatio,
+ taskDurationSeconds,
+ maxProcessingRatePerTask,
+ maxProcessingRatePerTask * 1.0
+ );
+ final int optimalTaskCount =
autoscaleSimulator.computeOptimalTaskCountInternal(costMetrics, true);
+ rows[i] = Map.of("lag", observedAggregateLag, "taskCount",
optimalTaskCount);
+ }
+
+ // Collect the results and return
+ return Map.of("data", rows);
+ }
+
/**
* Stops a supervisor with a given id and then removes it from the list.
* <p/>
diff --git
a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/supervisor/SupervisorResource.java
b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/supervisor/SupervisorResource.java
index 10297e25a88..73ab668f16d 100644
---
a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/supervisor/SupervisorResource.java
+++
b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/supervisor/SupervisorResource.java
@@ -39,6 +39,7 @@ import org.apache.druid.error.DruidException;
import org.apache.druid.indexing.overlord.DataSourceMetadata;
import org.apache.druid.indexing.overlord.TaskMaster;
import
org.apache.druid.indexing.overlord.http.security.SupervisorResourceFilter;
+import
org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostBasedAutoScalerConfig;
import org.apache.druid.java.util.common.StringUtils;
import org.apache.druid.java.util.common.UOE;
import org.apache.druid.segment.incremental.ParseExceptionReport;
@@ -183,7 +184,9 @@ public class SupervisorResource
);
}
- /** Audits supervisor spec submissions that changed or restarted the
supervisor. */
+ /**
+ * Audits supervisor spec submissions that changed or restarted the
supervisor.
+ */
private void auditSupervisorUpdate(final SupervisorSpec spec, final
HttpServletRequest req)
{
final String auditPayload
@@ -534,6 +537,31 @@ public class SupervisorResource
);
}
+ @POST
+ @Path("/{id}/autoscaler/simulate")
+ @Consumes(MediaType.APPLICATION_JSON)
+ @Produces(MediaType.APPLICATION_JSON)
+ @ResourceFilters(SupervisorResourceFilter.class)
+ public Response simulateAutoscaling(
+ @PathParam("id") String supervisorId,
+ CostBasedAutoScalerConfig autoScalerConfig,
+ @QueryParam("maxProcessingRatePerTask") int maxProcessingRatePerTask,
+ @QueryParam("currentTaskCount") Integer currentTaskCount,
+ @Context HttpServletRequest request
+ )
+ {
+ return asLeaderWithSupervisorManager(
+ manager -> Response.ok(
+ manager.simulateAutoscaling(
+ supervisorId,
+ autoScalerConfig,
+ maxProcessingRatePerTask,
+ currentTaskCount
+ )
+ ).build()
+ );
+ }
+
@GET
@Path("/history")
@Produces(MediaType.APPLICATION_JSON)
@@ -562,7 +590,12 @@ public class SupervisorResource
{
if (count != null && count <= 0) {
return Response.status(Response.Status.BAD_REQUEST)
- .entity(ImmutableMap.of("error",
StringUtils.format("Count must be greater than zero if set (count was %d)",
count)))
+ .entity(ImmutableMap.of("error",
+ StringUtils.format(
+ "Count must be greater than
zero if set (count was %d)",
+ count
+ )
+ ))
.build();
}
diff --git
a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java
b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java
index d9f84c5b703..6efd24d4b09 100644
---
a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java
+++
b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java
@@ -1657,7 +1657,7 @@ public abstract class
SeekableStreamSupervisor<PartitionIdType, SequenceOffsetTy
boolean includeOffsets
)
{
- int numPartitions =
partitionGroups.values().stream().mapToInt(Set::size).sum();
+ final int numPartitions = getKnownPartitionCount();
final SeekableStreamSupervisorReportPayload<PartitionIdType,
SequenceOffsetType> payload = createReportPayload(
numPartitions,
@@ -3195,6 +3195,21 @@ public abstract class
SeekableStreamSupervisor<PartitionIdType, SequenceOffsetTy
return false;
}
+ /**
+ * The number of partitions as last fetched from the underlying stream.
+ * This method differs from {@link #getPartitionCount()} as it does not
+ * refetch the current partition count from the stream, thus avoiding the
+ * need for locks or a network call.
+ */
+ public int getKnownPartitionCount()
+ {
+ return partitionIds.size();
+ }
+
+ /**
+ * Fetches the current partition count from the underlying stream using the
+ * {@link #recordSupplier}.
+ */
public int getPartitionCount()
{
recordSupplierLock.lock();
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 b1955f8428f..455ecb8e099 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
@@ -63,6 +63,7 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
{
private static final EmittingLogger log = new
EmittingLogger(CostBasedAutoScaler.class);
+ public static final String AUTOSCALER_TYPE_NAME = "costBased";
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";
@@ -98,6 +99,7 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
private final String supervisorId;
private final SeekableStreamSupervisor supervisor;
private final ServiceEmitter emitter;
+ private final boolean isSimulator;
private final SupervisorSpec spec;
private final CostBasedAutoScalerConfig config;
private final ScheduledExecutorService autoscalerExecutor;
@@ -123,10 +125,34 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
this.costFunction = new WeightedCostFunction();
this.autoscalerExecutor =
Execs.scheduledSingleThreaded("CostBasedAutoScaler-"
+
StringUtils.encodeForFormat(spec.getId()));
+ this.isSimulator = false;
+ }
+
+ private CostBasedAutoScaler(CostBasedAutoScalerConfig config, String
supervisorId)
+ {
+ this.config = config;
+ this.costFunction = new WeightedCostFunction();
+ this.isSimulator = true;
+ this.supervisorId = "simulator__" + supervisorId;
+
+ this.spec = null;
+ this.supervisor = null;
+ this.emitter = null;
+ this.processingRateSamples = null;
+ this.autoscalerExecutor = null;
+ }
+
+ public static CostBasedAutoScaler createSimulator(CostBasedAutoScalerConfig
config, String supervisorId)
+ {
+ return new CostBasedAutoScaler(config, supervisorId);
}
private ServiceMetricEvent.Builder getMetricBuilder()
{
+ if (isSimulator) {
+ return ServiceMetricEvent.builder();
+ }
+
return
ServiceMetricEvent.builder()
.setDimension(DruidMetrics.SUPERVISOR_ID,
supervisorId)
@@ -236,12 +262,24 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
* metrics are unusable. Returning the current task count means the current
count is already
* optimal (or no better candidate could be evaluated).
*/
- int computeOptimalTaskCount(CostMetrics metrics)
+ public int computeOptimalTaskCount(CostMetrics metrics)
+ {
+ return computeOptimalTaskCountInternal(metrics, false);
+ }
+
+ /**
+ * 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
+ * optimal (or no better candidate could be evaluated).
+ *
+ * @param isSimulation enables or disables task count computation logs and
metrics
+ */
+ public int computeOptimalTaskCountInternal(CostMetrics metrics, boolean
isSimulation)
{
final Either<String, Boolean> result = validateMetricsForScaling(metrics);
if (result.isError()) {
log.debug("Valid metrics are not yet available for scaling
supervisor[%s]", supervisorId);
- emitter.emit(
+ emitMetric(
getMetricBuilder()
.setDimension(DruidMetrics.DESCRIPTION, result.error())
.setMetric(INVALID_METRICS_COUNT, 1L)
@@ -263,21 +301,23 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
if (validTaskCounts.length == 0) {
// Return current count (not an error) so the supervisor can clamp it
back into bounds.
- log.warn("No valid task counts after applying constraints for
supervisor[%s]", supervisorId);
+ if (!isSimulation) {
+ log.warn("No valid task counts after applying constraints for
supervisor[%s]", supervisorId);
+ }
return currentTaskCount;
}
final boolean highLag = isHighLag(metrics);
final boolean criticalLag = isCriticalLag(metrics);
if (criticalLag) {
- log.info(
+ logInfo(
"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_THRESHOLD_FRACTION * 100,
config.getCriticalLagThreshold()
);
} else if (highLag) {
- log.info(
+ logInfo(
"Supervisor[%s] aggregateLag[%.0f] crossed [%.0f%%] of
criticalLagThreshold[%d]: widening scale-up"
+ " candidates and maxing out the high-lag cost factor.",
supervisorId, metrics.getAggregateLag(),
WeightedCostFunction.HIGH_LAG_THRESHOLD_FRACTION * 100,
@@ -291,7 +331,7 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
CostResult optimalCost = currentCost;
final double idleRatioEstimatedFromRate =
metrics.estimateIdleRatioFromProcessingRate();
- log.info(
+ logInfo(
"Computing optimal taskCount for supervisor[%s] with metrics:"
+ " currentTaskCount[%d], avgPartitionLag[%.1f],
avgProcessingRate[%.1f], maxProcessingRate[%.1f]"
+ " idleRatio[%.1f], pollIdleRatio[%.1f], lagWeight[%.2f],
idleWeight[%.2f].",
@@ -345,43 +385,45 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
}
}
- emitter.emit(getMetricBuilder().setMetric(OPTIMAL_TASK_COUNT_METRIC,
(long) optimalTaskCount));
- emitter.emit(getMetricBuilder().setMetric(LAG_WEIGHT_METRIC,
config.getLagWeight()));
- 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) {
- emitter.emit(getMetricBuilder().setMetric(AVG_PROCESSING_RATE_METRIC,
metrics.getAvgProcessingRate()));
- }
- if (metrics.getPollIdleRatio() >= 0) {
- emitter.emit(getMetricBuilder().setMetric(AVG_POLL_IDLE_RATIO,
metrics.getPollIdleRatio()));
- }
- if (idleRatioEstimatedFromRate >= 0) {
-
emitter.emit(getMetricBuilder().setMetric(IDLE_RATIO_ESTIMATED_FROM_RATE,
idleRatioEstimatedFromRate));
+ if (!isSimulation) {
+ emitMetric(getMetricBuilder().setMetric(OPTIMAL_TASK_COUNT_METRIC,
(long) optimalTaskCount));
+ emitMetric(getMetricBuilder().setMetric(LAG_WEIGHT_METRIC,
config.getLagWeight()));
+ emitMetric(getMetricBuilder().setMetric(IDLE_WEIGHT_METRIC,
config.getIdleWeight()));
+ emitMetric(getMetricBuilder().setMetric(CURRENT_LAG_COST_METRIC,
currentCost.lagCost()));
+ emitMetric(getMetricBuilder().setMetric(CURRENT_IDLE_COST_METRIC,
currentCost.idleCost()));
+ emitMetric(getMetricBuilder().setMetric(CURRENT_COST_METRIC,
currentCost.totalCost()));
+ emitMetric(getMetricBuilder().setMetric(OPTIMAL_COST_METRIC,
optimalCost.totalCost()));
+
+ // Emit avg rate and idle metrics only if they are available
+ if (metrics.getAvgProcessingRate() >= 0) {
+ emitMetric(getMetricBuilder().setMetric(AVG_PROCESSING_RATE_METRIC,
metrics.getAvgProcessingRate()));
+ }
+ if (metrics.getPollIdleRatio() >= 0) {
+ emitMetric(getMetricBuilder().setMetric(AVG_POLL_IDLE_RATIO,
metrics.getPollIdleRatio()));
+ }
+ if (idleRatioEstimatedFromRate >= 0) {
+
emitMetric(getMetricBuilder().setMetric(IDLE_RATIO_ESTIMATED_FROM_RATE,
idleRatioEstimatedFromRate));
+ }
+
+ if (optimalTaskCount != currentTaskCount) {
+ logInfo(
+ "Optimal taskCount[%d] for supervisor[%s] has lowest cost[%.4f]
out of the following candidates: %n%s",
+ optimalTaskCount, supervisorId, optimalCost.totalCost(),
constructCostTable(validTaskCounts, costResults)
+ );
+ emitMetric(getMetricBuilder().setMetric(OPTIMAL_LAG_COST_METRIC,
optimalCost.lagCost()));
+ emitMetric(getMetricBuilder().setMetric(OPTIMAL_IDLE_COST_METRIC,
optimalCost.idleCost()));
+ }
}
- if (optimalTaskCount != currentTaskCount) {
- log.info(
- "Optimal taskCount[%d] for supervisor[%s] has lowest cost[%.4f] out
of the following candidates: %n%s",
- optimalTaskCount, supervisorId, optimalCost.totalCost(),
constructCostTable(validTaskCounts, costResults)
- );
- emitter.emit(getMetricBuilder().setMetric(OPTIMAL_LAG_COST_METRIC,
optimalCost.lagCost()));
- emitter.emit(getMetricBuilder().setMetric(OPTIMAL_IDLE_COST_METRIC,
optimalCost.idleCost()));
-
- if (!criticalLag) {
- 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;
- }
+ if (!criticalLag) {
+ final double costDropPercent
+ = 100.0 * (currentCost.totalCost() - optimalCost.totalCost()) /
currentCost.totalCost();
+ if (costDropPercent < config.getMinCostDropPercentForScaling()) {
+ logInfo(
+ "Skipping scaling since cost drop percent[%.2f] is less than
required minCostDropPercentForScaling[%d]",
+ costDropPercent, config.getMinCostDropPercentForScaling()
+ );
+ return currentTaskCount;
}
}
@@ -603,4 +645,20 @@ public class CostBasedAutoScaler implements
SupervisorTaskAutoScaler
}
}
+ /**
+ * Emits metric for the given builder if this is not a simulator.
+ */
+ private void emitMetric(ServiceMetricEvent.Builder eventBuilder)
+ {
+ if (!isSimulator) {
+ emitter.emit(eventBuilder);
+ }
+ }
+
+ private void logInfo(String msgFormat, Object... args)
+ {
+ if (!isSimulator) {
+ log.info(msgFormat, args);
+ }
+ }
}
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 92dda7e775e..008df34c7a9 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
@@ -133,10 +133,11 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
if (this.enableTaskAutoScaler) {
Preconditions.checkNotNull(taskCountMax, "taskCountMax is required when
enableTaskAutoScaler is true");
Preconditions.checkNotNull(taskCountMin, "taskCountMin is required when
enableTaskAutoScaler is true");
+ Preconditions.checkArgument(taskCountMin >= 1, "taskCountMin must be at
least 1");
Preconditions.checkArgument(taskCountMax >= taskCountMin, "taskCountMax
must be >= taskCountMin");
Preconditions.checkArgument(
taskCountStart == null || (taskCountStart >= taskCountMin &&
taskCountStart <= taskCountMax),
- "taskCountMin <= taskCountStart <= taskCountMax"
+ "1 <= taskCountMin <= taskCountStart <= taskCountMax"
);
this.taskCountMax = taskCountMax;
this.taskCountMin = taskCountMin;
@@ -176,6 +177,31 @@ public class CostBasedAutoScalerConfig implements
AutoScalerConfig
return new Builder();
}
+ /**
+ * Config used to simulate the cost function without running an actual
supervisor.
+ */
+ public static CostBasedAutoScalerConfig forSimulation(
+ int taskCountMin,
+ int taskCountMax,
+ double optimalTaskIdleRatio,
+ @Nullable Double lagWeight,
+ @Nullable Double idleWeight
+ )
+ {
+ final Builder builder = builder()
+ .taskCountMin(taskCountMin)
+ .taskCountMax(taskCountMax)
+ .optimalTaskIdleRatio(optimalTaskIdleRatio)
+ .enableTaskAutoScaler(true);
+ if (lagWeight != null) {
+ builder.lagWeight(lagWeight);
+ }
+ if (idleWeight != null) {
+ builder.idleWeight(idleWeight);
+ }
+ return builder.build();
+ }
+
@Override
@JsonProperty
public boolean getEnableTaskAutoScaler()
diff --git
a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManagerTest.java
b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManagerTest.java
index bcc2e12e248..7fe9a0ad80b 100644
---
a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManagerTest.java
+++
b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManagerTest.java
@@ -46,6 +46,7 @@ import
org.apache.druid.indexing.seekablestream.supervisor.SeekableStreamSupervi
import
org.apache.druid.indexing.seekablestream.supervisor.SeekableStreamSupervisorIngestionSpec;
import
org.apache.druid.indexing.seekablestream.supervisor.SeekableStreamSupervisorSpec;
import
org.apache.druid.indexing.seekablestream.supervisor.SupervisorIOConfigBuilder;
+import
org.apache.druid.indexing.seekablestream.supervisor.autoscaler.CostBasedAutoScalerConfig;
import org.apache.druid.jackson.DefaultObjectMapper;
import org.apache.druid.java.util.common.DateTimes;
import org.apache.druid.java.util.common.Intervals;
@@ -626,6 +627,80 @@ public class SupervisorManagerTest extends EasyMockSupport
return (ConcurrentHashMap<String, Pair<Supervisor, SupervisorSpec>>)
field.get(manager);
}
+ @Test
+ public void testSimulateAutoscalingUsesLiveTaskCountAboveConfiguredMaximum()
throws Exception
+ {
+ final String supervisorId = "supervisor";
+ final SeekableStreamSupervisor<Integer, String, ByteEntity> supervisor =
EasyMock.createMock(
+ SeekableStreamSupervisor.class
+ );
+ final int partitionCount = 10;
+
EasyMock.expect(supervisor.getKnownPartitionCount()).andReturn(partitionCount);
+ EasyMock.replay(supervisor);
+
+ final TestBackfillSupervisorSpec.IngestionSpec ingestionSpec = new
TestBackfillSupervisorSpec.IngestionSpec(
+ new TestBackfillSupervisorSpec.IOConfig("test-stream", null, null)
+ );
+ getSupervisorsMap().put(
+ supervisorId,
+ Pair.of(supervisor, new TestBackfillSupervisorSpec(supervisorId,
ingestionSpec))
+ );
+
+ final Map<String, Object> result = manager.simulateAutoscaling(
+ supervisorId,
+ CostBasedAutoScalerConfig.forSimulation(1, 10, 0.1, null, null),
+ 100,
+ null
+ );
+
+ final Object[] data = (Object[]) result.get("data");
+ Assert.assertEquals(200, data.length);
+ Assert.assertTrue(data[0] instanceof Map);
+ final Map<?, ?> firstDataPoint = (Map<?, ?>) data[0];
+ Assert.assertTrue(firstDataPoint.containsKey("lag"));
+ Assert.assertTrue(firstDataPoint.containsKey("taskCount"));
+ for (Object dataPoint : data) {
+ Assert.assertTrue(dataPoint instanceof Map);
+ final Number taskCount = (Number) ((Map<?, ?>)
dataPoint).get("taskCount");
+ Assert.assertTrue(taskCount.intValue() >= 1 && taskCount.intValue() <=
10);
+ }
+ EasyMock.verify(supervisor);
+ }
+
+ @Test
+ public void
testSimulateAutoscalingRejectsExplicitTaskCountAboveConfiguredMaximum() throws
Exception
+ {
+ final String supervisorId = "supervisor";
+ final SeekableStreamSupervisor<Integer, String, ByteEntity> supervisor =
EasyMock.createMock(
+ SeekableStreamSupervisor.class
+ );
+ EasyMock.replay(supervisor);
+
+ final TestBackfillSupervisorSpec.IngestionSpec ingestionSpec = new
TestBackfillSupervisorSpec.IngestionSpec(
+ new TestBackfillSupervisorSpec.IOConfig("test-stream", null, null)
+ );
+ getSupervisorsMap().put(
+ supervisorId,
+ Pair.of(supervisor, new TestBackfillSupervisorSpec(supervisorId,
ingestionSpec))
+ );
+
+ MatcherAssert.assertThat(
+ Assert.assertThrows(
+ DruidException.class,
+ () -> manager.simulateAutoscaling(
+ supervisorId,
+ CostBasedAutoScalerConfig.forSimulation(1, 10, 0.1, null,
null),
+ 100,
+ 11
+ )
+ ),
+ DruidExceptionMatcher.invalidInput().expectMessageIs(
+ "Value of currentTaskCount[11] must be within taskCountMin[1] and
taskCountMax[10]"
+ )
+ );
+ EasyMock.verify(supervisor);
+ }
+
@Test
public void testHandoffTaskGroupsEarly()
{
diff --git
a/web-console/src/dialogs/supervisor-table-action-dialog/auto-scaler-panel/auto-scaler-panel.scss
b/web-console/src/dialogs/supervisor-table-action-dialog/auto-scaler-panel/auto-scaler-panel.scss
new file mode 100644
index 00000000000..add3b0902a0
--- /dev/null
+++
b/web-console/src/dialogs/supervisor-table-action-dialog/auto-scaler-panel/auto-scaler-panel.scss
@@ -0,0 +1,68 @@
+/*
+ * 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.
+ */
+
+.auto-scaler-panel {
+ display: flex;
+ flex-direction: column;
+ height: 100%;
+ padding: 15px;
+ overflow: hidden;
+
+ .auto-scaler-controls {
+ display: flex;
+ flex-wrap: wrap;
+ gap: 5px 30px;
+ margin-bottom: 15px;
+
+ .bp4-form-group {
+ margin: 0;
+ min-width: 220px;
+
+ .bp4-label {
+ white-space: nowrap;
+ }
+
+ .bp4-numeric-input {
+ width: 100px;
+ }
+
+ .bp4-slider {
+ width: 200px;
+ min-width: 200px;
+ margin-top: -15px;
+ }
+ }
+ }
+
+ .auto-scaler-chart-area {
+ position: relative;
+ flex: 1;
+ min-height: 300px;
+
+ .auto-scaler-echart {
+ width: 100%;
+ height: 100%;
+ min-height: 300px;
+ }
+
+ .auto-scaler-error {
+ color: #d5100a;
+ padding: 10px 0;
+ }
+ }
+}
diff --git
a/web-console/src/dialogs/supervisor-table-action-dialog/auto-scaler-panel/auto-scaler-panel.spec.tsx
b/web-console/src/dialogs/supervisor-table-action-dialog/auto-scaler-panel/auto-scaler-panel.spec.tsx
new file mode 100644
index 00000000000..efe649c68d9
--- /dev/null
+++
b/web-console/src/dialogs/supervisor-table-action-dialog/auto-scaler-panel/auto-scaler-panel.spec.tsx
@@ -0,0 +1,41 @@
+/*
+ * 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.
+ */
+
+import { getAutoScalerValidationError } from './auto-scaler-panel';
+
+describe('getAutoScalerValidationError', () => {
+ const validValues = {
+ taskCountMin: 1,
+ taskCountMax: 10,
+ maxProcessingRatePerTask: 1000,
+ optimalTaskIdleRatio: 0.2,
+ criticalLag: 1000000,
+ currentTaskCount: undefined,
+ };
+
+ it('validates the simulator inputs before making a request', () => {
+ expect(getAutoScalerValidationError(validValues)).toBeUndefined();
+ expect(getAutoScalerValidationError({ ...validValues, taskCountMin: 11
})).toBeDefined();
+ expect(
+ getAutoScalerValidationError({ ...validValues, maxProcessingRatePerTask:
99 }),
+ ).toBeDefined();
+ expect(getAutoScalerValidationError({ ...validValues,
optimalTaskIdleRatio: 0 })).toBeDefined();
+ expect(getAutoScalerValidationError({ ...validValues, criticalLag: 999
})).toBeDefined();
+ expect(getAutoScalerValidationError({ ...validValues, currentTaskCount: 11
})).toBeDefined();
+ });
+});
diff --git
a/web-console/src/dialogs/supervisor-table-action-dialog/auto-scaler-panel/auto-scaler-panel.tsx
b/web-console/src/dialogs/supervisor-table-action-dialog/auto-scaler-panel/auto-scaler-panel.tsx
new file mode 100644
index 00000000000..592c011761e
--- /dev/null
+++
b/web-console/src/dialogs/supervisor-table-action-dialog/auto-scaler-panel/auto-scaler-panel.tsx
@@ -0,0 +1,304 @@
+/*
+ * 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.
+ */
+
+import { FormGroup, NumericInput, Slider } from '@blueprintjs/core';
+import type { ECharts } from 'echarts';
+import * as echarts from 'echarts';
+import React, { useEffect, useMemo, useRef, useState } from 'react';
+
+import { Loader } from '../../../components/loader/loader';
+import { useQueryManager } from '../../../hooks';
+import { Api } from '../../../singletons';
+
+import './auto-scaler-panel.scss';
+
+interface AutoScalerRow {
+ lag: number;
+ taskCount: number;
+}
+
+interface AutoScalerPanelProps {
+ supervisorId: string;
+}
+
+export function getAutoScalerValidationError({
+ taskCountMin,
+ taskCountMax,
+ maxProcessingRatePerTask,
+ optimalTaskIdleRatio,
+ criticalLag,
+ currentTaskCount,
+}: {
+ taskCountMin: number;
+ taskCountMax: number;
+ maxProcessingRatePerTask: number;
+ optimalTaskIdleRatio: number;
+ criticalLag: number;
+ currentTaskCount: number | undefined;
+}): string | undefined {
+ if (taskCountMin > taskCountMax) return 'Minimum task count must not exceed
maximum task count';
+ if (maxProcessingRatePerTask < 100) return 'Max processing rate / task must
be at least 100';
+ if (optimalTaskIdleRatio <= 0 || optimalTaskIdleRatio >= 1) {
+ return 'Optimal task idle ratio must be greater than 0 and less than 1';
+ }
+ if (criticalLag < 1000) return 'Critical lag must be at least 1000';
+ if (
+ currentTaskCount !== undefined &&
+ (currentTaskCount < taskCountMin || currentTaskCount > taskCountMax)
+ ) {
+ return 'Current task count must be within the minimum and maximum task
count';
+ }
+ return undefined;
+}
+
+export const AutoScalerPanel = React.memo(function AutoScalerPanel(props:
AutoScalerPanelProps) {
+ const { supervisorId } = props;
+
+ const [taskCountMin, setTaskCountMin] = useState<number>(1);
+ const [taskCountMax, setTaskCountMax] = useState<number>(10);
+ const [maxProcessingRatePerTask, setMaxProcessingRatePerTask] =
useState<number>(1000);
+ const [optimalTaskIdleRatio, setOptimalTaskIdleRatio] =
useState<number>(0.2);
+ const [lagWeight, setLagWeight] = useState<number>(0.4);
+ // Idle weight is the complement of lag weight; one slider drives both.
+ const idleWeight = Math.round((1 - lagWeight) * 10) / 10;
+ const [criticalLag, setCriticalLag] = useState<number>(1000000);
+ // Undefined means "let the server use the supervisor's live task count".
+ const [currentTaskCount, setCurrentTaskCount] = useState<number |
undefined>(undefined);
+
+ const chartContainerRef = useRef<HTMLDivElement | undefined>(undefined);
+ const chartRef = useRef<ECharts | undefined>(undefined);
+ const query = useMemo(
+ () => ({
+ supervisorId,
+ taskCountMin,
+ taskCountMax,
+ maxProcessingRatePerTask,
+ optimalTaskIdleRatio,
+ lagWeight,
+ idleWeight,
+ criticalLag,
+ currentTaskCount,
+ }),
+ [
+ supervisorId,
+ taskCountMin,
+ taskCountMax,
+ maxProcessingRatePerTask,
+ optimalTaskIdleRatio,
+ lagWeight,
+ idleWeight,
+ criticalLag,
+ currentTaskCount,
+ ],
+ );
+ const validationError = getAutoScalerValidationError(query);
+
+ const [dataState] = useQueryManager<
+ {
+ supervisorId: string;
+ taskCountMin: number;
+ taskCountMax: number;
+ maxProcessingRatePerTask: number;
+ optimalTaskIdleRatio: number;
+ lagWeight: number;
+ idleWeight: number;
+ criticalLag: number;
+ currentTaskCount: number | undefined;
+ },
+ AutoScalerRow[]
+ >({
+ query: validationError ? undefined : query,
+ debounceIdle: 300,
+ debounceLoading: 500,
+ processQuery: async (params, signal) => {
+ const resp = await Api.instance.post<{ data: AutoScalerRow[] }>(
+
`/druid/indexer/v1/supervisor/${Api.encodePath(params.supervisorId)}/autoscaler/simulate`,
+ {
+ autoScalerStrategy: 'costBased',
+ enableTaskAutoScaler: true,
+ taskCountMin: params.taskCountMin,
+ taskCountMax: params.taskCountMax,
+ optimalTaskIdleRatio: params.optimalTaskIdleRatio,
+ lagWeight: params.lagWeight,
+ idleWeight: params.idleWeight,
+ criticalLagThreshold: params.criticalLag,
+ },
+ {
+ params: {
+ maxProcessingRatePerTask: params.maxProcessingRatePerTask,
+ currentTaskCount: params.currentTaskCount,
+ },
+ signal,
+ },
+ );
+ return resp.data.data ?? (resp.data as any);
+ },
+ });
+
+ function setupChart(container: HTMLDivElement): ECharts {
+ const myChart = echarts.init(container, 'dark');
+ myChart.setOption({
+ tooltip: {
+ trigger: 'axis',
+ },
+ grid: {
+ left: '3%',
+ right: '4%',
+ bottom: '3%',
+ containLabel: true,
+ },
+ xAxis: {
+ type: 'value',
+ name: 'Lag (records)',
+ nameLocation: 'middle',
+ nameGap: 30,
+ },
+ yAxis: {
+ type: 'value',
+ name: 'Task count',
+ nameLocation: 'middle',
+ nameGap: 40,
+ },
+ series: [
+ {
+ name: 'Task count',
+ type: 'line',
+ showSymbol: false,
+ data: [],
+ },
+ ],
+ });
+ return myChart;
+ }
+
+ useEffect(() => {
+ return () => {
+ chartRef.current?.dispose();
+ };
+ }, []);
+
+ useEffect(() => {
+ const myChart = chartRef.current;
+ const data = dataState.data;
+ if (!myChart || !data) return;
+
+ myChart.setOption({
+ series: [
+ {
+ data: data.map(row => [row.lag, row.taskCount]),
+ },
+ ],
+ });
+ }, [dataState.data]);
+
+ useEffect(() => {
+ const myChart = chartRef.current;
+ if (!myChart) return;
+ myChart.resize();
+ }, []);
+
+ const errorMessage = validationError ?? dataState.getErrorMessage();
+
+ return (
+ <div className="auto-scaler-panel">
+ <div className="auto-scaler-controls">
+ <FormGroup label="Min task count" inline>
+ <NumericInput
+ value={taskCountMin}
+ min={1}
+ max={taskCountMax}
+ onValueChange={v => setTaskCountMin(v)}
+ buttonPosition="none"
+ fill
+ />
+ </FormGroup>
+ <FormGroup label="Max task count" inline>
+ <NumericInput
+ value={taskCountMax}
+ min={taskCountMin}
+ onValueChange={v => setTaskCountMax(v)}
+ buttonPosition="none"
+ fill
+ />
+ </FormGroup>
+ <FormGroup label="Max processing rate / task" inline>
+ <NumericInput
+ value={maxProcessingRatePerTask}
+ min={100}
+ onValueChange={v => setMaxProcessingRatePerTask(v)}
+ buttonPosition="none"
+ fill
+ />
+ </FormGroup>
+ <FormGroup label="Optimal task idle ratio" inline>
+ <NumericInput
+ value={optimalTaskIdleRatio}
+ min={0.01}
+ max={0.99}
+ stepSize={0.1}
+ minorStepSize={0.01}
+ onValueChange={v => setOptimalTaskIdleRatio(v)}
+ fill
+ />
+ </FormGroup>
+ <FormGroup label="Current task count" inline>
+ <NumericInput
+ value={currentTaskCount ?? ''}
+ min={taskCountMin}
+ max={taskCountMax}
+ placeholder="Supervisor's current"
+ onValueChange={v => setCurrentTaskCount(isNaN(v) ? undefined : v)}
+ buttonPosition="none"
+ fill
+ />
+ </FormGroup>
+ <FormGroup label="Critical lag (records)" inline>
+ <NumericInput
+ value={criticalLag}
+ min={1000}
+ onValueChange={v => setCriticalLag(v)}
+ buttonPosition="none"
+ fill
+ />
+ </FormGroup>
+ <FormGroup label={`Lag ${lagWeight.toFixed(1)} / Idle
${idleWeight.toFixed(1)} weight`}>
+ <Slider
+ min={0}
+ max={1}
+ stepSize={0.1}
+ labelStepSize={0.5}
+ value={lagWeight}
+ onChange={v => setLagWeight(Math.round(v * 10) / 10)}
+ />
+ </FormGroup>
+ </div>
+ <div className="auto-scaler-chart-area">
+ {errorMessage && <div
className="auto-scaler-error">{errorMessage}</div>}
+ {dataState.loading && <Loader />}
+ <div
+ className="auto-scaler-echart"
+ ref={container => {
+ if (chartRef.current || !container) return;
+ chartContainerRef.current = container;
+ chartRef.current = setupChart(container);
+ }}
+ />
+ </div>
+ </div>
+ );
+});
diff --git
a/web-console/src/dialogs/supervisor-table-action-dialog/supervisor-table-action-dialog.tsx
b/web-console/src/dialogs/supervisor-table-action-dialog/supervisor-table-action-dialog.tsx
index 44ef05d5d59..edea317e035 100644
---
a/web-console/src/dialogs/supervisor-table-action-dialog/supervisor-table-action-dialog.tsx
+++
b/web-console/src/dialogs/supervisor-table-action-dialog/supervisor-table-action-dialog.tsx
@@ -26,12 +26,14 @@ import type { BasicAction } from '../../utils/basic-action';
import type { SideButtonMetaData } from
'../table-action-dialog/table-action-dialog';
import { TableActionDialog } from '../table-action-dialog/table-action-dialog';
+import { AutoScalerPanel } from './auto-scaler-panel/auto-scaler-panel';
import { SupervisorStatisticsTable } from
'./supervisor-statistics-table/supervisor-statistics-table';
-type SupervisorTableActionDialogTab = 'status' | 'stats' | 'spec' | 'history';
+type SupervisorTableActionDialogTab = 'status' | 'stats' | 'spec' | 'history'
| 'auto-scaler';
interface SupervisorTableActionDialogProps {
supervisorId: string;
+ supervisorType?: string;
actions: BasicAction[];
onClose: () => void;
}
@@ -39,9 +41,11 @@ interface SupervisorTableActionDialogProps {
export const SupervisorTableActionDialog = React.memo(function
SupervisorTableActionDialog(
props: SupervisorTableActionDialogProps,
) {
- const { supervisorId, actions, onClose } = props;
+ const { supervisorId, supervisorType, actions, onClose } = props;
const [activeTab, setActiveTab] =
useState<SupervisorTableActionDialogTab>('status');
+ const isKafka = supervisorType === 'kafka';
+
const supervisorTableSideButtonMetadata: SideButtonMetaData[] = [
{
icon: 'dashboard',
@@ -67,6 +71,16 @@ export const SupervisorTableActionDialog =
React.memo(function SupervisorTableAc
active: activeTab === 'history',
onClick: () => setActiveTab('history'),
},
+ ...(isKafka
+ ? [
+ {
+ icon: 'predictive-analysis' as const,
+ text: 'Simulate auto-scaler',
+ active: activeTab === 'auto-scaler',
+ onClick: () => setActiveTab('auto-scaler'),
+ },
+ ]
+ : []),
];
const supervisorEndpointBase =
`/druid/indexer/v1/supervisor/${Api.encodePath(supervisorId)}`;
@@ -98,6 +112,7 @@ export const SupervisorTableActionDialog =
React.memo(function SupervisorTableAc
/>
)}
{activeTab === 'history' && <SupervisorHistoryPanel
supervisorId={supervisorId} />}
+ {activeTab === 'auto-scaler' && isKafka && <AutoScalerPanel
supervisorId={supervisorId} />}
</TableActionDialog>
);
});
diff --git a/web-console/src/views/supervisors-view/supervisors-view.tsx
b/web-console/src/views/supervisors-view/supervisors-view.tsx
index 632853e6126..88ff023a75a 100644
--- a/web-console/src/views/supervisors-view/supervisors-view.tsx
+++ b/web-console/src/views/supervisors-view/supervisors-view.tsx
@@ -212,6 +212,7 @@ export interface SupervisorsViewState {
alertErrorMsg?: string;
supervisorTableActionDialogId?: string;
+ supervisorTableActionDialogType?: string;
supervisorTableActionDialogActions: BasicAction[];
visibleColumns: LocalStorageBackedVisibility;
@@ -810,6 +811,7 @@ export class SupervisorsView extends React.PureComponent<
private onSupervisorDetail(supervisor: SupervisorQueryResultRow) {
this.setState({
supervisorTableActionDialogId: supervisor.supervisor_id,
+ supervisorTableActionDialogType: supervisor.type,
supervisorTableActionDialogActions:
this.getSupervisorActions(supervisor),
});
}
@@ -1326,6 +1328,7 @@ export class SupervisorsView extends React.PureComponent<
supervisorSpecDialogOpen,
alertErrorMsg,
supervisorTableActionDialogId,
+ supervisorTableActionDialogType,
supervisorTableActionDialogActions,
visibleColumns,
} = this.state;
@@ -1393,8 +1396,14 @@ export class SupervisorsView extends React.PureComponent<
{supervisorTableActionDialogId && (
<SupervisorTableActionDialog
supervisorId={supervisorTableActionDialogId}
+ supervisorType={supervisorTableActionDialogType}
actions={supervisorTableActionDialogActions}
- onClose={() => this.setState({ supervisorTableActionDialogId:
undefined })}
+ onClose={() =>
+ this.setState({
+ supervisorTableActionDialogId: undefined,
+ supervisorTableActionDialogType: undefined,
+ })
+ }
/>
)}
</div>
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]