[ 
https://issues.apache.org/jira/browse/FLINK-40407?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18122184#comment-18122184
 ] 

Kaustubh Butte commented on FLINK-40407:
----------------------------------------

I opened a pull request for this ticket:
[https://github.com/apache/flink-kubernetes-operator/pull/1219] 

It is stacked on the pull request for FLINK-40405, because both change
AutoscalerFlinkMetrics, ScalingExecutor and JobAutoScalerImpl. Please review
FLINK-40405 first. I will rebase this one after that one merges.

The change adds a ScaleReason enum with four values: backlog, high_load,
input_spike and low_util. JobVertexScaler gets the causes from the same values
that make the target capacity. A vertex reports a set, because a backlog and a
high load can apply to one decision. The causes of all the scaled vertices make
one sorted tag, for example "backlog|high_load". Four causes give a maximum of
15 tag values.

The untagged autoscaler.scalings counter keeps its total.
ScalingSummary.scaleReasons is @JsonIgnore, so the state store needs no
migration.

The change is in the flink-autoscaler module. A standalone or YARN deployment
that uses the library also gets the tags.

Please assign the ticket to me.

> [FLINK-autoscaler] Add reason tags to `autoscaler.scalings` counter to 
> explain why scaling was triggered
> --------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40407
>                 URL: https://issues.apache.org/jira/browse/FLINK-40407
>             Project: Flink
>          Issue Type: Improvement
>          Components: Autoscaler
>    Affects Versions: 2.0.0, 1.20.1
>            Reporter: Kaustubh Butte
>            Priority: Major
>              Labels: autoscaling, observability, pull-request-available
>
> Currently, the `autoscaler.scalings` counter is incremented whenever a 
> scaling decision is applied, but there is no metric indicating WHY the 
> autoscaler decided to scale. Understanding the scaling trigger is critical 
> for:
>  - Incident triage ("did it scale because of a traffic spike or a Kafka lag 
> buildup?")
>  - Validating autoscaler behavior ("is this job scaling due to genuine load 
> or a false-positive signal?")
>  - Building targeted alerts ("alert on backlog-driven scale-ups but not 
> routine load adjustments")
> The information needed to classify the reason is already computed in the 
> autoscaler's evaluation pipeline but is discarded after the scaling decision 
> is made.
> *Proposed reasons:*
> *Scale-up:*
> |Reason|Condition|Meaning|
> |`backlog`|`CATCH_UP_DATA_RATE.current > 0` and contributes >10% of 
> `targetCapacity`|Consumer lag accumulated; scaling up to catch up|
> |`high_load`|`CATCH_UP_DATA_RATE ~ 0` and `TRUE_PROCESSING_RATE < 
> SCALE_UP_RATE_THRESHOLD`|Sustained input rate exceeds processing capacity|
> |`input_spike`|`TARGET_DATA_RATE.current > 1.5 * 
> TARGET_DATA_RATE.average`|Sudden traffic increase|
> *Scale-down:*
> |Reason|Condition|Meaning|
> |`low_utilization`|`TRUE_PROCESSING_RATE > 
> SCALE_DOWN_RATE_THRESHOLD`|Processing capacity far exceeds demand|
> *Classification logic:*
> The target capacity is already decomposed in 
> `AutoScalerUtils.getTargetProcessingCapacity()`:
> ```
> targetCapacity = lagCatchupRate + restartCatchupRate + 
> inputTargetAtUtilization
> ```
> A single vertex can have multiple reasons simultaneously (e.g., backlog AND 
> high load both contributing to targetCapacity). And a single scaling cycle 
> evaluates all vertices, each potentially with different reasons. Therefore, 
> reasons are collected as a set per vertex, then aggregated across all 
> vertices into a sorted pipe-delimited string.
> *Per-vertex classification* (inside 
> `JobVertexScaler.computeScaleTargetParallelism()`, inline with the decision):
>  - If `lagCatchupRate / targetCapacity > 0.1` then add `backlog`
>  - If `TARGET_DATA_RATE.current / TARGET_DATA_RATE.average > 1.5` then add 
> `input_spike`
>  - If scale-up and no specific reason matched, add `high_load` (fallback)
>  - If scale-down, add `low_util`
> *Job-level aggregation* (in `ScalingExecutor`):
>  - Collect all reason sets from all vertices in the scaling summary
>  - Merge into a single `TreeSet<String>` (alphabetically sorted)
>  - Join with `|` delimiter: e.g., `"backlog|high_load"`
> This approach prevents algorithm-reason drift because the classification is 
> done inline in the same function that computes the scale factor, reading the 
> same computed variables. If the algorithm changes, the developer sees the 
> reason-tagging code right there.
> *Maximum cardinality:* 15 distinct tag values (2^4 - 1 combinations of 4 
> atomic reasons). In practice fewer, since `high_load` is only added when no 
> other scale-up reason applies.
> *Implementation:*
> 1. Add `Set<String> scaleReasons` to `ParallelismChange` in 
> `JobVertexScaler`, tag reasons inline alongside scale factor computation
> 2. Add `Set<String> reasons` to `ScalingSummary`, propagated from 
> `ParallelismChange`
> 3. In `ScalingExecutor`, aggregate reasons from all vertices into a sorted 
> pipe-delimited string
> 4. In `AutoscalerFlinkMetrics`, register tagged sub-counters dynamically 
> using `MetricGroup.addGroup("reason", reasonTag).counter("scalings")`
> 5. Keep the existing `numScalings` counter incrementing for backward 
> compatibility
> *Result in metric reporters:*
>  - `autoscaler.scalings` (total, backward compatible)
>  - `autoscaler.scalings \{reason=backlog}` – pure backlog catch-up
>  - `autoscaler.scalings \{reason=backlog|high_load}` – backlog + steady-state 
> overload
>  - `autoscaler.scalings \{reason=high_load}` – sustained input rate exceeds 
> capacity
>  - `autoscaler.scalings \{reason=input_spike}` – sudden traffic increase
>  - `autoscaler.scalings \{reason=low_util}` – scale-down, overcapacity
>  - (and other sorted combinations up to 15 total)
> *Files affected:*
>  - `JobVertexScaler.java` – add `Set<String> scaleReasons` to 
> `ParallelismChange`, tag reasons inline
>  - `ScalingSummary.java` – add `Set<String> reasons` field
>  - `ScalingExecutor.java` – aggregate reasons across vertices into 
> pipe-delimited string
>  - `AutoscalerFlinkMetrics.java` – add dynamically registered tagged counters
>  - `JobAutoScalerImpl.java` – pass aggregated reason string to metrics
> All changes are in the `flink-autoscaler` module (standalone, 
> platform-agnostic). Works with both K8s and YARN deployments.
> *Relationship to other tickets:*
> This is a companion to the `balanced` counter disambiguation ticket. Both can 
> be implemented together or independently. Together they provide complete 
> observability over every autoscaler decision cycle: why it scaled, why it 
> didn't scale, or why it errored.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to