kaustubhbutte17 opened a new pull request, #1219:
URL: https://github.com/apache/flink-kubernetes-operator/pull/1219
## What is the purpose of the change
Stacked on #1218. Please review that one first. This pull request therefore
shows two commits against `main`, and will be rebased once #1218 merges.
The `autoscaler.scalings` counter reports that an evaluation cycle changed
the
parallelism. It does not report what drove the change. A scale up caused by a
Kafka backlog, a scale up caused by a sustained load, and a scale down
caused by
spare capacity all look the same.
This matters for two reasons. An operator who watches a rollout cannot tell a
healthy reaction to traffic from an autoscaler that fights a downstream
bottleneck. And an alert on scaling activity has to fire on all of it,
because
there is no way to select only the backlog driven changes.
This change tags the counter with what drove the decision.
## Brief change log
- Add `ScaleReason`, an enum with four values: `backlog`, `high_load`,
`input_spike` and `low_util`.
- `JobVertexScaler.computeScaleReasons` names the reasons from the same
values
that `AutoScalerUtils.getTargetProcessingCapacity` uses to build the target
capacity.
- `ParallelismChange` and `ScalingSummary` carry the reasons of a vertex.
- `ScalingExecutor` unions the reasons over every scaled vertex and publishes
them on the `ScalingCycleState`, once the decision is final.
- `AutoscalerFlinkMetrics.incrementScaling(Set<ScaleReason>)` increments the
existing untagged counter, then a counter under the `reason` metric
variable.
- Update the metric and the internals documentation, in English and Chinese.
## Design notes
**The reasons sit next to the decision.** `computeScaleReasons` reads
`CATCH_UP_DATA_RATE` and `TARGET_DATA_RATE`, which are the same terms the
target
capacity formula sums. A reason derived in another class drifts away from the
algorithm as the algorithm changes, and the metric then lies.
**A set, not one value.** The target capacity is the lag catch up rate plus
the
restart catch up rate plus the rate at the target utilization. A backlog and
a
high load can both contribute to one decision, so a vertex reports every
reason
that applies. The reasons of all scaled vertices join into one sorted tag,
for
example `backlog|high_load`.
**Cardinality.** Four reasons give at most 15 tag values, that is 2^4 - 1. A
counter is created the first time its combination occurs, so a job carries
only
the combinations it really produces.
**Backward compatibility.** The untagged `autoscaler.scalings` counter keeps
its
meaning and its total.
**No state migration.** `ScalingSummary.scaleReasons` is `@JsonIgnore`. The
scaling history in the state store keeps the shape it had before, so an
operator
can roll back to an earlier operator version without touching the state.
**`INPUT_SPIKE_FACTOR` is a constant, not a config option.** It is 1.5. The
tag
only labels a decision and never changes one, so a configuration key would
add
surface for no behavioural gain. Happy to make it configurable if reviewers
prefer.
**The reasons stay out of `equals` and `hashCode`**, for the same reason as
`NoChangeReason` in the parent commit.
`testScaleReasonIsExcludedFromEquality` guards this.
## Verifying this change
`mvn -pl flink-autoscaler test` passes, with 314 tests.
New tests:
- `JobVertexScalerTest.testScaleReasonIsLowUtilOnScaleDown`
- `JobVertexScalerTest.testScaleReasonIsHighLoadOnScaleUpWithoutBacklog`
- `JobVertexScalerTest.testScaleReasonReportsBacklogAndHighLoadTogether`
- `JobVertexScalerTest.testScaleReasonReportsAnInputSpike`
- `JobVertexScalerTest.testScaleReasonTagIsSortedAndStable`
- `JobVertexScalerTest.testScaleReasonIsExcludedFromEquality`
- `AutoScalerFlinkMetricsTest.testScalingsCounterIsTaggedWithTheReasons`
-
`AutoScalerFlinkMetricsTest.testScalingsCounterIsNotRegisteredWithoutAReason`
Each new test was confirmed to fail when the behaviour it covers is broken on
purpose, and to pass again after the code is restored.
## Does this pull request potentially affect one of the following parts
- Dependencies (does it add or upgrade a dependency): **no**
- The public API, i.e., is any changed class annotated with
`@Public(Evolving)`: **no**
- Core observability classes: **yes**, `AutoscalerFlinkMetrics` gains
counters.
Existing counters are unchanged.
- The serialized state format: **no**, the new field is not serialized.
- The autoscaler decision logic: **no**, every scaling decision is identical.
Only the reporting of the decision changes.
## Documentation
- Does this pull request introduce a new feature? **yes**, per reason
counters
on `autoscaler.scalings`.
- If yes, how is the feature documented? **docs**, in
`docs/content/docs/operations/metrics.md` and
`docs/content/docs/internals/autoscaler.md`, with the Chinese mirrors
updated.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]