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