dev-donghwan opened a new pull request, #4566: URL: https://github.com/apache/flink-cdc/pull/4566
## What is the purpose of this pull request? Jira: https://issues.apache.org/jira/browse/FLINK-40853 ### Problem With `scan.lsn-commit.checkpoints-num-delay` (default 3), `PostgresSourceReader` commits the offset of an older checkpoint to the replication slot. A checkpoint has no stream offset in `lastCheckpointOffsets` when it was taken: * before the stream split was assigned, e.g. during the snapshot phase of the `initial` startup mode, or * while newly added tables were being snapshotted with `scan.newly-added-table.enabled` (`PostgresSourceReader#snapshotState` removes its offset). Because of the delay, such a checkpoint is still committed after the stream phase has started. `PostgresStreamFetchTask#commitCurrentOffset(null)` then falls back to the `LAST_COMMIT_LSN` of the live Debezium offset context, which can be beyond the offset of the latest completed checkpoint. PostgreSQL does not resend changes before `confirmed_flush_lsn`, so after a failover the changes between the checkpoint offset and the committed LSN are lost for good. #2539 switched this method to commit the checkpoint's LSN for exactly this reason, and its comment describes the scenario. The live LSN was only kept as a fallback for a `null` offset, which FLINK-35129 made reachable by adding the commit delay. ### Changes * `PostgresStreamFetchTask#commitCurrentOffset`: skip the commit when the checkpoint has no stream offset, instead of committing the live LSN. The slot then only receives offsets of completed checkpoints, and catches up at the next checkpoint that has one. * The now unused read of the live offset context is removed. `commitCurrentOffset` -> `PostgresStreamingChangeEventSource#commitOffset` -> `ReplicationStream#flushLsn` is the only path that moves the slot. ### Tests New tests in `PostgresSourceReaderTest`, running against the PostgreSQL test container. All three fail on master and pass with this change: * `testLsnCommitDoesNotPassLatestCompletedCheckpoint`: checkpoints 1-3 are taken before the stream split is assigned, checkpoint 4 after change 3001 in the stream phase, then change 3002 is read and checkpoints 1-7 complete. On master, the slot moves to `0/17B9A70` while checkpoint 4 is at `0/17B98F0`. * `testNoDataLossWhenRestoringAfterLsnCommit`: the same, then restores from checkpoint 4 and inserts 3003. On master only 3003 arrives. 3002 is lost. * `testLsnCommitDoesNotPassLatestCompletedCheckpointAfterNewlyAddedTables`: offset commit is turned off with `OffsetCommitEvent` for checkpoints 4-6 and on again for checkpoint 7, then change 3002 is read and checkpoint 7 completes. On master, the slot moves to `0/17B9AA8` while checkpoint 7 is at `0/17B9928`. Existing tests pass locally: `PostgresSourceReaderTest`, `NewlyAddedTableITCase` and `PostgresSourceITCase` (78 tests, 1 skipped by an existing assumption). -- 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]
