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

ASF GitHub Bot updated FLINK-40853:
-----------------------------------
    Labels: pull-request-available  (was: )

> [postgres] Data loss: LSN commit passes the latest completed checkpoint when 
> the delayed checkpoint has no stream offset
> ------------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40853
>                 URL: https://issues.apache.org/jira/browse/FLINK-40853
>             Project: Flink
>          Issue Type: Bug
>          Components: Flink CDC
>    Affects Versions: cdc-3.2.0, cdc-3.6.0
>            Reporter: kim
>            Priority: Critical
>              Labels: pull-request-available
>
> h3. Summary
> With {{scan.lsn-commit.checkpoints-num-delay}} (default 3), the Postgres 
> incremental source commits the offset of an older checkpoint to the 
> replication slot. If that checkpoint has no stream offset, 
> {{PostgresStreamFetchTask#commitCurrentOffset(null)}} falls back to the 
> {{LAST_COMMIT_LSN}} of the live Debezium offset context. That LSN can be 
> ahead of the offset of the latest completed checkpoint. PostgreSQL never 
> resends changes before {{confirmed_flush_lsn}}, so after restoring from that 
> checkpoint the transactions in between are lost for good. Nothing is logged 
> or reported.
> h3. When it happens
> A checkpoint has no stream offset in 
> {{IncrementalSourceReaderWithCommit#lastCheckpointOffsets}} in two cases. 
> Because of the commit delay, such a checkpoint is still committed after the 
> stream phase has started.
> # The checkpoint was taken before the stream split was assigned, for example 
> during the snapshot phase of the {{initial}} startup mode. The first 
> checkpoints after switching to the stream phase (3 with the default delay) 
> commit the live LSN.
> # The checkpoint was taken while newly added tables were being snapshotted 
> with {{scan.newly-added-table.enabled}}. 
> {{PostgresSourceReader#snapshotState}} removes its offset. Once offset commit 
> is turned back on, the delayed checkpoints commit the live LSN.
> A failover in these windows loses the changes between the checkpoint offset 
> and the committed LSN.
> h3. Cause
> * [#2539|https://github.com/apache/flink-cdc/pull/2539] (cdc-3.1.0) changed 
> {{commitCurrentOffset}} to commit the checkpoint's LSN instead of the live 
> one, because committing the live LSN loses data after a failover. Its comment 
> explains exactly this scenario. It kept the live LSN only as a fallback when 
> the offset is {{null}}.
> * FLINK-35129 (cdc-3.2.0) added the commit delay. The delay can pick a 
> checkpoint without a stream offset, so the {{null}} fallback is actually used 
> and brings back the data loss #2539 fixed.
> The only path that moves the slot is {{commitCurrentOffset}} -> 
> {{PostgresStreamingChangeEventSource#commitOffset}} -> 
> {{ReplicationStream#flushLsn}}.
> h3. Reproduction
> Tests added to {{PostgresSourceReaderTest}}, using the PostgreSQL test 
> container (PostgreSQL 14):
> * {{testLsnCommitDoesNotPassLatestCompletedCheckpoint}}: checkpoints 1-3 are 
> taken before the stream split is assigned, checkpoint 4 in the stream phase 
> after change 3001, then change 3002 is read and checkpoints 1-7 complete. On 
> master, the offset of checkpoint 4 is {{0/17B98F0}}, but the slot moves to 
> {{0/17B9A70}}.
> * {{testNoDataLossWhenRestoringAfterLsnCommit}}: same as above, then restores 
> from checkpoint 4 and inserts 3003. On master only 3003 arrives. 3002 is lost.
> * {{testLsnCommitDoesNotPassLatestCompletedCheckpointAfterNewlyAddedTables}}: 
> offset commit is turned off by {{OffsetCommitEvent}} for checkpoints 4-6 and 
> on again for checkpoint 7, then change 3002 is read and checkpoint 7 
> completes. On master, the offset of checkpoint 7 is {{0/17B9928}}, but the 
> slot moves to {{0/17B9AA8}}.
> All three fail on master and pass with the fix below. We also hit the first 
> case with a real job (initial startup, TaskManager killed right after the 
> switch to the stream phase), which lost a transaction every time in 4 runs.
> This is a client-side issue: Flink CDC reports an LSN it has not 
> checkpointed. It does not depend on the PostgreSQL version.
> h3. Workaround
> Set {{scan.lsn-commit.checkpoints-num-delay}} to 0.
> h3. Proposed fix
> Skip the commit in {{commitCurrentOffset}} when the checkpoint has no stream 
> offset. The slot then only receives offsets of completed checkpoints, and 
> catches up at the next checkpoint that has one.
> h3. Not covered
> FLINK-39265 is a different data loss after recovery, caused by 
> {{WalPositionLocator}} filtering when reading. It is independent of this 
> issue: with the fix, the restore test above receives 3002 through the same 
> recovery path.



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

Reply via email to