kim created FLINK-40853:
---------------------------

             Summary: [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


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 (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