li3zhi4 opened a new issue, #11676: URL: https://github.com/apache/seatunnel/issues/11676
### Search before asking - [X] I had searched in the [issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue+label%3A%22bug%22) and found no issue describing the CDC **enumerator** restored-split re-assignment race: - #10574 "[Bug][connector-jdbc] Data loss when sink fails and job restores in exactly-once mode" — same user-visible symptom (sink fails, job restores, no more data), but maintainers confirmed the root cause is the **Kafka source** restore-offset logic, fixed by PR #10612; the CDC `IncrementalSourceEnumerator` path is untouched. - #4773 / #4595 (2023, both closed/stale) — restore-path issues in a different direction (reader not yet registered when split is assigned; `RestoredSplitOperation` NPE); neither covers this race. ### What happened When a MySQL CDC job fails (e.g. the sink hits a transient error), SeaTunnel restarts the job from the last successful checkpoint. The checkpoint continues to complete, the job stays `RUNNING`, but the MySQL CDC source **never writes any data again** — the incremental split is never re-assigned to the reader. Restart logs show the CDC enumerator restored with empty split state: ```text SnapshotSplitAssigner created with remaining tables: [] SnapshotSplitAssigner created with remaining splits: [] SnapshotSplitAssigner created with assigned splits: [] ``` and no log of a reader receiving the restored split. Every subsequent checkpoint is an empty checkpoint. ### SeaTunnel Version 2.3.13 (verified locally). **Also reproducible on current upstream `dev`** (verified 2026-08-06 with a regression test, see below). ### SeaTunnel Config ```hocon source { MySQL-CDC { parallelism = 1 url = "jdbc:mysql://<host>:3306/<db>" username = "user" password = "pass" database-names = ["<db>"] table-names = ["<db>.<table>"] } } sink { Doris { ... } # any sink whose transient failure triggers engine-side restart } ``` ### Running Command ```shell ./bin/seatunnel.sh --config ./config/xxx.conf ``` ### Error Exception No error. The task restarts, checkpoints complete, but the CDC source produces no records — a silent data stall. ### Root cause (confirmed by code inspection + reproduction on dev) Reader-side checkpoint-restored splits are handed back to the CDC enumerator through `RestoredSplitOperation`, which calls `IncrementalSourceEnumerator.addSplitsBack()`. The original implementation only calls `splitAssigner.addSplits(splits)`: ```java // IncrementalSourceEnumerator.java (current dev, original) @Override public void addSplitsBack(List<SourceSplitBase> splits, int subtaskId) { LOG.debug("Incremental Source Enumerator adds splits back: {}", splits); splitAssigner.addSplits(splits); // ❌ no re-run of assignSplits() } ``` If the restored split arrives **after** the reader has already sent its split request (the reader is parked in `readersAwaitingSplit`), the split sits in the assigner but `assignSplits()` is never re-invoked, so the reader never receives the split and the job idles forever. `handleSplitRequest()` does invoke `assignSplits()` when the request arrives — but only at that moment; a split restored later never triggers another assignment. ### Proposed fix (implemented & verified locally on 2.3.13 and on dev) ```java @Override public synchronized void addSplitsBack(List<SourceSplitBase> splits, int subtaskId) { LOG.debug("Incremental Source Enumerator adds splits back: {}", splits); splitAssigner.addSplits(splits); if (running) { assignSplits(); // ✅ hand restored splits to the currently waiting reader } } ``` - `IncrementalSourceEnumerator.addSplitsBack()` becomes synchronous: after adding restored splits, if the enumerator is already running, immediately run `assignSplits()` so the restored split is dispatched to the waiting reader and consumption resumes from the checkpointed offset. - Non-goals: does not touch checkpoint semantics, sink stream-load transactions, or the transient sink error itself. ### Verification 1. **Unit test (deterministic)**: `IncrementalSourceEnumeratorTest#shouldAssignRestoredSplitsToWaitingReader` — reader already waiting, restored split must be dispatched immediately. 2. **Reproduced on upstream dev (2026-08-06)**: the unit test **fails on the original code** (`expected: <[SnapshotSplit(tableId=null, ...)]> but was: <[]>` — restored split never assigned) and **passes after the 3-line fix** (`Tests run: 1, Failures: 0`). 3. **Official E2E (Zeta + Testcontainers)**: `MysqlCDCIT#testMysqlCdcContinuesAfterSinkFailure` — completes snapshot/checkpoint, injects a sink failure, waits for Zeta auto-restart, removes the fault, writes new binlog data, and asserts source/sink consistency again (1/1 passed, ~95 s). ### Zeta or Flink or Spark Version zeta ### Java or Scala Version Java 8 / 11 (both OK) ### Screenshots _No response_ ### Are you willing to submit PR? - [X] Yes I am willing to submit a PR! (fix implemented locally on 2.3.13 and ported to `dev`) ### Code of Conduct - [X] I agree to follow this project's [Code of Conduct](https://www.apache.org/foundation/policies/conduct) -- 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]
