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]

Reply via email to