wuhainan opened a new issue, #12348: URL: https://github.com/apache/seatunnel/issues/12348
### 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 similar issues. ### What happened A MySQL-CDC multi-table job was restored from a savepoint after adding a new table, `alpha_online.accounts`. The job configuration uses `startup.mode = initial`. The new table is discovered during restore, and the JobManager logs: ```log SnapshotSplitAssigner created with remaining tables: [alpha_online.accounts] ``` However, the TaskManager restores an existing `incremental-split-0` containing only the previously configured tables. It does not contain `alpha_online.accounts`. No snapshot split for `alpha_online.accounts` is assigned to or executed by the TaskManager. Therefore, the target table is created but receives no historical records. After restore, the long-running incremental split and Debezium schema history still do not include `alpha_online.accounts`. When MySQL emits a binlog event for this new table, the source task fails with: ```log io.debezium.DebeziumException: Encountered change event for table alpha_online.accounts whose schema isn't known to this connector ``` The failure propagates through the CDC fetcher and causes the Flink job to restart repeatedly: ```log java.lang.RuntimeException: One or more fetchers have encountered exception Caused by: org.apache.kafka.connect.errors.ConnectException: An exception occurred in the change event producer. This connector will be stopped. ``` With a failure-rate restart strategy configured as 10 failures within 300 seconds, the job eventually terminates with: ```log Recovery is suppressed by FailureRateRestartBackoffTimeStrategy ``` The issue is therefore not only that the new table does not receive its initial snapshot. Adding a table and restoring from an existing savepoint also makes the existing CDC job unavailable as soon as binlog events for the new table arrive. ## Expected behavior When a new table is added to a CDC job restored from a savepoint/checkpoint, SeaTunnel should initialize the new table's schema, execute its initial snapshot, and include it in subsequent incremental CDC processing while preserving the existing tables' restored offsets. ## Possible fix The restore logic should treat a newly added table as a state migration, not only add it to `SnapshotPhaseState.remainingTables`. A possible solution is: 1. During restore, detect tables present in the current configuration but absent from the restored source state. 2. Initialize Debezium schema metadata for each new table before consuming any binlog event for that table. 3. Schedule and assign snapshot splits for new tables even when a restored reader already owns a long-running `IncrementalSplit`. 4. After the new table snapshot finishes, atomically update or recreate the incremental split so that: - its `tableIds` includes the new table; - it contains the new table's snapshot watermark / offset information; - its schema and history state includes the new table. 5. Preserve the restored offsets for existing tables. The new table must start from the snapshot-consistent binlog watermark, so that neither data loss nor duplicate events are introduced. Please add integration tests for both savepoint and checkpoint recovery: - Start a MySQL-CDC job with table A. - Complete at least one checkpoint/savepoint. - Add table B to the job configuration. - Restore from the previous state. - Verify that B receives an initial snapshot. - Insert/update rows in B after restore. - Verify that B continues to receive incremental changes. - Verify that A continues from its original restored offset without duplicate snapshots. ### SeaTunnel Version 2.3.13 ### SeaTunnel Config ```conf env { parallelism = 1 job.mode = "STREAMING" job.name = "mysql2tidb_au_omnibus" checkpoint.interval = 30000 checkpoint.timeout = 600000 checkpoint.mode = "EXACTLY_ONCE" restart-strategy = "failure-rate" restart-strategy.failure-rate.max-failures-per-interval = 10 restart-strategy.failure-rate.failure-rate-interval = "300 s" restart-strategy.failure-rate.delay = "10 s" } source { MySQL-CDC { hostname = "<redacted>" port = 3306 username = "<redacted>" password = "<redacted>" server-id = 6333 server-time-zone = "Asia/Shanghai" startup.mode = "initial" database-names = ["alpha_online"] # Existing tables were already present when the savepoint was created. # alpha_online.accounts was added after the savepoint was created. table-names = [ "alpha_online.sub_account_applications", "alpha_online.sub_account_profiles", "alpha_online.deposit_application_plans", "alpha_online.withdrawals", "alpha_online.deposit_notices", "alpha_online.position_transfer_requests", "alpha_online.deposit_applications", "alpha_online.account_profiles", "alpha_online.accounts" ] } } sink { # JDBC / MultiTableSink configuration omitted because it is unrelated # to the source failure. All credentials and addresses are redacted. } ``` ### Running Command ```shell # The job is submitted by an internal platform using Flink YARN Application mode. # The equivalent command is: $FLINK_HOME/bin/flink run-application \ -t yarn-application \ -c org.apache.seatunnel.core.starter.flink.SeaTunnelFlink \ /path/to/seatunnel-flink-20-starter.jar \ --config /path/to/v2.conf ``` ### Error Exception ```log java.lang.RuntimeException: One or more fetchers have encountered exception at org.apache.seatunnel.connectors.seatunnel.common.source.reader.fetcher.SplitFetcherManager.checkErrors(SplitFetcherManager.java:147) ... Caused by: java.lang.RuntimeException: SplitFetcher thread 0 received unexpected exception while polling the records ... Caused by: org.apache.kafka.connect.errors.ConnectException: An exception occurred in the change event producer. This connector will be stopped. ... Caused by: io.debezium.DebeziumException: Error processing binlog event ... Caused by: io.debezium.DebeziumException: Encountered change event for table alpha_online.accounts whose schema isn't known to this connector at io.debezium.connector.mysql.MySqlStreamingChangeEventSource.informAboutUnknownTableIfRequired(MySqlStreamingChangeEventSource.java:768) ... The source task fails repeatedly after restoring from the savepoint. Eventually Flink terminates the job: org.apache.flink.runtime.JobException: Recovery is suppressed by FailureRateRestartBackoffTimeStrategy( failuresIntervalMS=300000, backoffTimeMS=10000, maxFailuresPerInterval=10 ) ``` ### Zeta or Flink or Spark Version flink 1.20 ### Java or Scala Version java 17 ### Screenshots _No response_ ### Are you willing to submit PR? - [x] Yes I am willing to submit a PR! ### 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]
