DanielLeens commented on PR #12258:
URL: https://github.com/apache/seatunnel/pull/12258#issuecomment-5697629769
Root cause found and fixed in `bd08491a32` (production code,
`connector-cdc-mariadb` only; one file, no shared code touched).
**Evidence from this head's own CI run** (`merlau/seatunnel` run
`34848994932`, `updated-modules-integration-test-part-6 (8)`): the diagnostic
dump shows the sink at `[1, 2, 3, 11, 12, 12, 12, 21, 22]` vs. source `[1, 2,
3, 11, 12, 12, 21, 22]` - exactly one extra copy of `12`, the last row
delivered before the restored checkpoint. The restored source state logged by
the restore job (`Restore create source from checkpoint state`) is:
```
startupOffset={transaction_id=null, ts_sec=1789397780,
file=mariadb-bin.000002, pos=92893, gtids=1-223344-57, row=1, server_id=223344,
event=3}
```
That is a mid-transaction *restart* offset - the Debezium `sourceOffset()`
attached to the last emitted row: the enclosing transaction's start position
plus the `event`/`row` counters of that transaction already emitted. Debezium
(`MySqlOffsetContext.Loader.load` -> `setInitialSkips`) re-reads the
transaction from `pos` and uses those two counters to skip what it already
delivered.
**The bug**: `MariaDbSourceFetchTaskContext.loadStartingOffsetState` forced
`event`/`row` to `0` for every startup mode other than `specific` before
handing the map to the loader (`MariaDbSourceFetchTaskContext.java:285-291`
before this fix). So on restore the transaction at `pos=92893` was replayed
with nothing skipped and its row (`12`) was emitted a second time.
`MySqlSourceFetchTaskContext` passes `offset.getOffset()` through untouched,
which is why the identical `MysqlCDCCheckpointRestoreIT` passes on `dev`. The
two other restore tests in this class pass because their restored offsets came
from watermark/heartbeat positions with no counters at all (e.g. the savepoint
test's restored state: `{..., pos=71176, gtids=1-223344-44,
server_id=223344}`), so there was nothing to lose.
**Fix**: hand the offset map to the Debezium loader untouched, exactly as
the MySQL connector does (offsets without counters are loaded as `0` by
`longOffsetValue`, so watermark/`initial`/`earliest`/`timestamp` startups are
unaffected; `specific` was never altered).
`MariaDbIncrementalSourceStartupConfigTest`/`MariaDbIncrementalSplitStateTest`
already assert that user-supplied and checkpointed skip counters are preserved,
which this now honours end to end. The regression coverage is
`testMysqlCdcRestoresAfterCheckpointedFullPipelineFailure` itself - it should
go green on this head; I have not observed that run yet.
--
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]