li3zhi4 opened a new issue, #11617:
URL: https://github.com/apache/seatunnel/issues/11617

   ### 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 similar issues: #6789 and #7400, but both were auto-closed by the 
stale bot without a fix.
   
   ### What happened
   
   When a MySQL CDC job configures `stop.mode = "specific"` with a specific 
binlog stop offset, the job **does not stop** at the specified position — it 
keeps running forever even after the stop offset has been passed.
   
   This is exactly the same symptom reported in #6789 and #7400, which were 
closed without a resolution. The root cause is confirmed by code inspection 
below and a fix has been implemented and verified locally on the 2.3.13 branch.
   
   ### SeaTunnel Version
   
   2.3.13 (also reproducible on 2.3.4 / 2.3.6 per #6789 / #7400)
   
   ### SeaTunnel Config
   
   ```conf
   source {
     MySQL-CDC {
       parallelism = 1
       url = "jdbc:mysql://<host>:3306/<db>"
       username = "user"
       password = "pass"
       database-names = ["<db>"]
       table-names = ["<db>.<table>"]
   
       startup.mode = "specific"
       startup.specific-offset.file = "mysql-bin.059734"
       startup.specific-offset.pos = 4
   
       stop.mode = "specific"
       stop.specific-offset.file = "mysql-bin.059818"
       stop.specific-offset.pos = 4
     }
   }
   
   sink {
     Console {}
   }
   ```
   
   ### Running Command
   
   ```shell
   ./bin/seatunnel.sh --config ./config/xxx.conf
   ```
   
   ### Error Exception
   
   ```log
   No error. The job simply does not end even after reaching the position 
designated as stop.mode.
   ```
   
   ### Root cause (confirmed by code inspection)
   
   SeaTunnel 2.3.13 has two binlog-reading classes:
   
   | Class | Used for | stopOffset check | Stops? |
   |---|---|---|---|
   | `MySqlBinlogSplitReadTask` | binlog backfill of **snapshot splits** | ✅ 
yes | ✅ yes |
   | `MySqlBinlogFetchTask` | binlog reading of the **incremental split** | ❌ 
**no** | ❌ **no** |
   
   ```java
   // MySqlBinlogFetchTask.java (2.3.13, original)
   @Override
   public void execute(FetchTask.Context context) throws Exception {
       if (startupMode.equals(StartupMode.TIMESTAMP)) {
           mySqlStreamingChangeEventSource = new 
TimestampFilterMySqlStreamingChangeEventSource(...);
       } else {
           // ❌ created directly, NO stopOffset check, runs forever
           mySqlStreamingChangeEventSource = new 
MySqlStreamingChangeEventSource(...);
       }
       mySqlStreamingChangeEventSource.execute(...); // infinite
   }
   ```
   
   `stop.mode = "specific"` is mainly used to make the **incremental split** 
stop, but the incremental split's reader `MySqlBinlogFetchTask` has **no 
stopOffset logic at all**. The upstream `b71d8739d5` only refactored the 
`startup.mode`/`stop.mode` options; it did not fix the binlog fetch behavior.
   
   Two related secondary bugs were also found and fixed locally:
   
   1. **`BinlogOffset.compareTo()` returns wrong result when GTID sets are 
mixed** — the comparison may fall back incorrectly when one side has a GTID 
set, breaking stop/start offset comparison.
   2. **`taskStarted` race condition in `IncrementalSourceStreamFetcher`** — a 
stale `taskStarted` flag could make the job finish prematurely (data silently 
missing, no error), affecting `stop.mode = "specific"` and `"latest"` (not 
`"never"`).
   
   ### Proposed fix (implemented & verified locally)
   
   - `MySqlBinlogFetchTask`: add a `BoundedMySqlStreamingChangeEventSource` 
inner class (~200 lines) that checks the stop offset per event and ends the 
binlog read at the target position; keep the unbounded path unchanged.
   - Base layer: `Offset.isNeverStop()` default + `IncrementalSourceReader` 
allows incremental splits to finish normally + `IncrementalSourceStreamFetcher` 
returns `null` to signal split completion only for bounded reads.
   - `BinlogOffset`: fix `compareTo()` GTID fallback and add `isNeverStop()`.
   - `IncrementalSourceStreamFetcher`: set `taskStarted`/`executing` 
synchronously to close the bounded-read race.
   
   Statistics: 6 files changed, ~+280 / ~-15 / ~-10 lines. Verified with 
`stop.mode = "specific"` on a production-like setup; the job now terminates at 
the exact binlog position.
   
   ### 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 already implemented on the 
custom branch; will port to `dev`)
   
   ### Code of 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