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]
