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)