[ 
https://issues.apache.org/jira/browse/FLINK-40113?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Piotr Nowojski closed FLINK-40113.
----------------------------------
    Fix Version/s: 2.4.0
       Resolution: Fixed

merged commit 
[{{47abc18}}|https://github.com/apache/flink/commit/47abc18f7ac249c193f72697e8dad448920936a0]
 into apache:master 

> State restore does not validate max-parallelism agreement between operators 
> sharing a vertex
> --------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40113
>                 URL: https://issues.apache.org/jira/browse/FLINK-40113
>             Project: Flink
>          Issue Type: Improvement
>          Components: Runtime / Checkpointing, Runtime / State Backends
>    Affects Versions: 2.0.0, 1.18.0, 1.19.0, 1.20.0, 2.1.0, 2.2.0, 2.3.0, 2.4.0
>            Reporter: Rion Williams
>            Assignee: Rion Williams
>            Priority: Major
>              Labels: pull-request-available
>             Fix For: 2.4.0
>
>
> h2. Summary
> [FLINK-31996|https://issues.apache.org/jira/projects/FLINK/issues/FLINK-31996]
>  allowed operators with different max parallelism to be chained into a single 
> vertex, and made it the default 
> ({{{}pipeline.operator-chaining.chain-operators-with-different-max-parallelism{}}}).
>  On restore, Flink validates that an operator's max parallelism matches the 
> savepoint's and rejects an incompatible change with a clear error but that 
> validation is per-operator against the vertex's single max parallelism. It 
> does not verify that the operators now sharing a vertex have compatible max 
> parallelism with each other. When they don't, the per-vertex reconciliation 
> silently adopts one operator's value instead of rejecting.
> h2. Current behavior
>  - A single operator whose max parallelism changed since the savepoint → 
> rejected: "{_}Max parallelism mismatch … changed in a non-compatible way."{_} 
> (Correct.)
>  - Multiple operators chained into one vertex with different recorded max 
> parallelism into {{StateAssignmentOperation#checkParallelismPreconditions}} 
> which reconciles each against the shared vertex and adopts one value; 
> operators whose recorded value differs are accepted under the adopted 
> key-group count with {*}no equivalent check{*}.
> h2. The gap
> {{checkParallelismPreconditions}} checks each operator against the vertex's 
> max parallelism (and {{canRescaleMaxParallelism}} correctly permits rescaling 
> an auto-derived value), but nothing checks that the operators sharing the 
> vertex agree with one another. 
> [FLINK-31996|https://issues.apache.org/jira/projects/FLINK/issues/FLINK-31996]
>  made this arrangement reachable by default; the restore-side validation was 
> not extended to cover it. The single-operator path is guarded; the 
> shared-vertex path is not.
> h2. Why it matters
> A keyed operator can then be restored under a key-group count different from 
> the one its state was written with, so its state is mapped through a 
> different {{hash % maxParallelism}} layout. The restore should reject this, 
> consistent with how it already rejects the single-operator case, rather than 
> proceeding.
> h2. Proposed change (validation-only, low-risk)
> Add explicit validation during per-vertex reconciliation in 
> {{StateAssignmentOperation#checkParallelismPreconditions}}: for any vertex 
> that carries keyed state, verify that all operators chained into it recorded 
> the same maximum parallelism in the checkpoint, and fail with a clear error 
> otherwise. The
> comparison spans all chained operators, not only the keyed ones, because a 
> non-keyed operator's recorded value can win the reconciliation and misroute a
> keyed operator's state.
> This extends the check 
> [FLINK-31996|https://issues.apache.org/jira/projects/FLINK/issues/FLINK-31996]
>  left uncovered and rejects a configuration the current reconciliation 
> handles nondeterministically: operators on a keyed vertex with different 
> recorded max parallelism may restore correctly or silently corrupt keyed 
> state depending on the order their values are applied. Restores where those 
> operators agree (the normal case) are unaffected.
> h2. Reproduction Steps
> Using parallelism = 1 and default configuration, with a plain {{keyBy}}:
>  # Build {{source -> keyBy -> keyed process -> stateful map -> sink}}, where 
> the keyed {{process}} is a chain head with an auto-derived max parallelism 
> (128 at parallelism 1) and the downstream stateful {{map}} has a different 
> explicit {{setMaxParallelism}} (e.g. 64 or 256).
>  # Run with operator chaining disabled, process records, and trigger a 
> canonical savepoint. The keyed operator and the downstream operator are 
> separate vertices, each recording its own max parallelism.
>  # Restore with chaining enabled (the default). The downstream operator now 
> chains under the keyed head, so both operators' differing recorded max 
> parallelism values are assigned to the single merged keyed vertex.
>  # {{checkParallelismPreconditions}} reconciles each recorded value onto the 
> shared vertex one at a time with no chain-head precedence, so whichever is 
> applied last wins. When the downstream operator's value wins, the keyed 
> operator is restored under a foreign key-group count and its state is mapped 
> through a different {{hash % maxParallelism}} layout, either failing the 
> restore outright or silently dropping per-key state depending on the 
> direction of the mismatch. The reconciliation order is not deterministic, so 
> the same job can range from a correct restore to silent loss across runs.
> A self-contained ITCase reproducing this on both the HashMap and RocksDB 
> state backends is available  
> [here|https://gist.github.com/rionmonster/ea8bcb8bbd22cb2921bdaec878cc41ca].
>  



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

Reply via email to