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]

Reply via email to