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]

Reply via email to