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]

Reply via email to