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

   ### Search before asking
   
   - [x] I had searched in the issues and found no similar issues.
   
   ### What happened
   
   Flink CDC schema evolution needs stronger ordering and recovery guarantees 
when schema changes are processed with parallel sinks, checkpoints, XA 
transactions, TaskManager failures, and rescaling.
   
   The existing coordinator-based approach introduces several correctness and 
maintainability concerns:
   
   - Schema coordination depends on process-local coordinator state and 
additional control protocols that are difficult to recover consistently across 
Flink versions.
   - The source-side protocol state must be restored atomically. Restoring 
these fields independently may split logically related state during rescaling.
   - A failure after an external DDL succeeds but before the corresponding 
Flink state is checkpointed can replay the DDL.
   - Existing TaskManager recovery coverage did not prove recovery together 
with a real JDBC XA sink and global committer.
   
   The improvement is to replace the custom schema coordinator with a 
checkpointed data-plane protocol:
   - Assign each schema change a stable (producerId, sequence) identifier.
   - Persist the complete source protocol state as one atomic Flink 
operator-state entry.
   - Allow only the sink subtask owning the table key to emit the physical DDL.
   - Buffer dependent rows until their required schema sequence has been 
applied.
   - Restore schema snapshots, sequences, pending controls, and pending rows 
from Flink-managed state.
   
   The Flink translation layer is responsible for ordering, dependency gating, 
checkpoint state, recovery, and replay. Physical DDL execution and convergence 
remain the responsibility of each connector. External DDL therefore follows 
at-least-once control delivery with connector-provided convergent application 
rather than transactional exactly-once DDL.
   
   ### SeaTunnel Version
   --
   
   ### SeaTunnel Config
   ```
   env {
       parallelism = 5
       job.mode = "STREAMING"
       checkpoint.interval = 30000
       flink.pipeline.max-parallelism = 128
       flink.pipeline.operator-chaining = false
       flink.execution.checkpointing.unaligned.enabled = true
   }
   
   source {
       MySQL-CDC {
         server-id = 5752-5760
         username = "<source-user>"
         password = "<source-password>"
         table-names = ["shop.products"]
         url = "jdbc:mysql://mysql:3306/shop"
         schema-changes.enabled = true
         incremental.parallelism = 1
       }
   }
   
   sink {
       jdbc {
         parallelism = 5
         url = "jdbc:mysql://mysql:3306/shop"
         driver = "com.mysql.cj.jdbc.Driver"
         user = "<sink-user>"
         password = "<sink-password>"
         generate_sink_sql = true
         database = shop
         table = products_sink
         primary_keys = ["id"]
         is_exactly_once = true
         xa_data_source_class_name = "com.mysql.cj.jdbc.MysqlXADataSource"
       }
   }
   ```
   
   ### Running Command
   --
   
   ### Error Exception
   --
   
   ### Zeta or Flink or Spark Version
   --
   
   ### Java or Scala Version
   --
   
   ### 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.


-- 
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