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

kim updated FLINK-40853:
------------------------
    Description: 
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.


  was:
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.



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