zhangshenghang opened a new issue, #12382: URL: https://github.com/apache/seatunnel/issues/12382
### Search before asking - [x] I searched in the [issues](https://github.com/apache/seatunnel/issues) and found nothing similar. ### What happened `PostgresCDCIT.testPostgresCdcSnapshotOnlyAndCommittedOffsetStartupModes` failed once on the dev schedule run with a **data loss after restore** signature: the restored CDC job stayed healthy, kept completing checkpoints, and its WAL reader committed change-log offsets **past** the inserted row's LSN — but the row never reached the sink within 180 seconds. ### Evidence (run 35244548330, job `all-connectors-it-7 (11, ubuntu-latest)`, 2026-09-17) Test = zeta only, single invocation, 17:38:35 → 17:42:55 (259s), failure: ``` ConditionTimeoutException: expected: <1> but was: <0> within 3 minutes (PostgresCDCIT.java:625) ``` The final stage restores job `4570933213989348932` from savepoint, waits for the replication slot to be re-attached, then inserts row `id=15` and waits up to 180s for it in the sink: | time (UTC) | event | |---|---| | 17:39:44 | `seatunnel.sh -s 4570933213989348932` (savepoint; checkpoint id 4) | | 17:39:51 | `pipeline(1) restore with SAVEPOINT on checkpointId(4)`, tasks start | | 17:39:52 | slot `seatunnel_12bb6b32e3e1d547` re-attached: `START_REPLICATION ... LOGICAL 0/22BC640` | | ~17:39:53+ | test inserts `id=15` (only after slot active, per test design) | | 17:39:51.999 | `PostgresWalFetchTask - Start streaming change event source ... lsn=36423232` (= 0/22BC640) | | 17:40:01 | `PostgresWalFetchTask - Committing offset LSN{0/22BC778}` | | 17:40:31 | `PostgresWalFetchTask - Committing offset LSN{0/22BE610}` | | 17:39:56→17:42:55 | checkpoints keep completing every 5s (checkpoint id 5..30+) | | 17:42:55 | await times out; sink has 0 rows with `id = 15` | `commitChangeLogOffset` is fed the offset of the **last emitted change event**, so `Committing offset LSN{0/22BC778}` (300+ bytes past the restore LSN, logged seconds after the insert) strongly suggests the WAL fetcher did decode and "emit" the `id=15` event, yet it never arrived at the sink. No job failure was reported — `assertJobHasNoAsyncFailure` inside the await never tripped. So somewhere along `PostgresWalFetchTask → IncrementalSource reader → JDBC sink writer` after a **savepoint/restore cycle**, a decoded event is dropped or the sink writer loses the restored write path, while offset bookkeeping and checkpoints continue as if the record was delivered. ### Expected behavior A row inserted after the restored job re-attaches its replication slot must appear in the sink. ### Suspect areas - `PostgresWalFetchTask` emit/suppression logic right after restore (startupOffset handling vs. records already buffered in the Debezium queue) - reader ↔ enumerator split state restored from checkpoint (`checkpointId(4)`) vs. the freshly created fetcher - JDBC sink writer state after restore (e.g. restored writer writing to a stale batch/connection) ### Zeta or Flink Zeta -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
