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

kim commented on FLINK-40216:
-----------------------------

I'd like to work on this by porting the FLINK-38218 fix to flink-cdc-base. 
Please assign it to me.

> [cdc-base] Meta groups of finished snapshot splits become inconsistent after 
> restore due to lexicographic re-sorting of assigned splits
> ---------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40216
>                 URL: https://issues.apache.org/jira/browse/FLINK-40216
>             Project: Flink
>          Issue Type: Bug
>          Components: Flink CDC
>    Affects Versions: cdc-3.1.0, cdc-3.6.0
>            Reporter: kim
>            Priority: Major
>
> h3. Problem
> In {{flink-cdc-base}}, the order of {{assignedSplits}} is not preserved 
> across a checkpoint/savepoint restore, so the meta groups of finished 
> snapshot splits are partitioned differently before and after a restart.
> During a normal run, splits are appended to {{assignedSplits}} (a 
> {{LinkedHashMap}}) in *assignment order* ({{SnapshotSplitAssigner#getNext}}). 
> On restore, however:
> # {{PendingSplitsStateSerializer#readAssignedSnapshotSplits}} deserializes 
> the checkpointed splits into a plain {{HashMap}}, losing the checkpointed 
> order.
> # The restore constructor of {{SnapshotSplitAssigner}} then re-sorts the 
> entries *lexicographically* by split id (introduced in FLINK-34634 as a 
> workaround for the lost order):
> {code:java}
> this.assignedSplits =
>         assignedSplits.entrySet().stream()
>                 .sorted(Map.Entry.comparingByKey())
>                 ...
> {code}
> Split ids are generated as {{tableId + ":" + chunkId}} 
> ({{SnapshotSplit#generateSplitId}}), so with 10+ chunks the lexicographic 
> order ({{t:0, t:1, t:10, t:11, t:2, ...}}) differs from the assignment order 
> ({{t:0, t:1, t:2, ..., t:11}}). The same happens across tables when the 
> alphabetical order of table names differs from the capture order.
> h3. Impact
> {{IncrementalSourceEnumerator#sendStreamMetaRequestEvent}} partitions 
> {{splitAssigner.getFinishedSplitInfos()}} by {{chunk-meta.group.size}} 
> ({{Lists.partition}}) and serves the group requested by the reader. A reader 
> whose stream split had synchronized only part of the meta groups before the 
> restart resumes requesting from its next group id, but the enumerator now 
> serves groups computed from a differently-ordered list.
> We reproduced this with a unit test against current master (single table, 12 
> chunks, {{chunk-meta.group.size = 2}}, reader synchronized meta groups 0..2 
> before the restart, group ids computed with 
> {{IncrementalSourceReader#getNextMetaGroupId}}): after the restore the reader 
> ends up with
> {code}
> [t:0, t:1, t:2, t:3, t:4, t:5, t:4, t:5, t:6, t:7, t:8, t:9]
> {code}
> i.e. {{t:4}}/{{t:5}} are received twice and {{t:10}}/{{t:11}} are never 
> delivered, while the stream split is considered complete because 
> {{StreamSplit#isCompletedSplit}} only compares sizes. Subsequently 
> {{IncrementalSourceStreamFetcher#shouldEmit}} finds no matching 
> {{FinishedSnapshotSplitInfo}} for change events in the key ranges of the lost 
> splits and silently drops them — permanent data loss in the stream phase. 
> Depending on where the restart happens, the mismatch can also surface as the 
> "invalid request meta group id" error.
> The MySQL connector carried the identical re-sort and was fixed by 
> FLINK-38218 (cdc-3.6.0), which removed it and made the whole pipeline 
> preserve the assignment order. As its author noted there: "Note that other 
> connectors that use SnapshotSplitAssigner from the flink-cdc-base package may 
> be prone to this issue as well." {{flink-cdc-base}} still carries the 
> identical code, affecting all connectors built on the base incremental 
> framework: Postgres, MongoDB, Oracle, SqlServer, Db2 (and their SQL/pipeline 
> wrappers). 
> {{MySqlSnapshotSplitAssignerTest#testFinishedSnapshotSplitInfosAreInOrderOfAssignment}}
>  guards this property on the MySQL side; {{flink-cdc-base}} has no equivalent 
> test.
> h3. Proposed fix
> Port FLINK-38218 to {{flink-cdc-base}}:
> * deserialize {{assignedSplits}} into a {{LinkedHashMap}} in 
> {{PendingSplitsStateSerializer}} (the serialized form already preserves 
> iteration order) and declare the order contract on the field, instead of 
> re-sorting in the restore constructor of {{SnapshotSplitAssigner}};
> * remove the lexicographic sort in {{HybridSplitAssigner#createStreamSplit}};
> * replace the order-dependent deduplication in {{IncrementalSourceReader}} 
> ({{getExistedSplitsOfLastGroup}}) with the order-agnostic prefix-discard 
> approach from FLINK-38218;
> * guard {{StreamSplit}} against duplicated split infos, mirroring 
> {{MySqlBinlogSplit#ensureNoDuplicates}}.
> This fix is independent of the metadata-transmission protocol redesign 
> discussed in FLINK-38270 and can be applied separately.



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

Reply via email to