CloverDew opened a new pull request, #11960:
URL: https://github.com/apache/seatunnel/pull/11960

   ### Purpose of this pull request
   Close: https://github.com/apache/seatunnel/issues/11958
   
   This PR replaces the process-local Flink schema coordinator and its 
coordinator protocol with a checkpointed data-plane schema-evolution protocol.
   
   The new protocol keeps schema controls and dependent data consistently 
ordered without introducing a custom OperatorCoordinator:
   
   - Stores the pending DDL, buffered rows, checkpoint fence, producer ID, 
sequence, and last emitted change ID in one atomic source protocol state.
   - Assigns one sink gate as the physical DDL owner for each table.
   - Buffers data rows that arrive before their required schema control.
   - Emits the schema command before dependent rows on the same downstream 
channel.
   - Checkpoints applied sequences, latest target schemas, out-of-order 
controls, and pending rows using Flink managed state.
   - Reconstructs the target schema before restored rows are sent to a newly 
created sink writer.
   
   The translation layer guarantees ordering, dependency gating, checkpoint 
state, and replay; Database DDL is delivered at least once relative to Flink 
checkpoints, Connectors remain responsible for making applySchemaChange 
synchronous and observable, propagating failures, and converging repeated DDL 
execution to the requested target schema.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Yes.
   
   Previously, Flink CDC schema evolution relied on a process-local coordinator 
and did not provide sufficient distributed recovery evidence for the physical 
DDL/checkpoint failure window.
   
   After this change:
   
   - Flink CDC schema evolution uses a checkpointed data-plane protocol.
   - Sink parallelism greater than one is supported by assigning one table 
owner and broadcasting compact schema controls.Schema and dependent data are 
restored and replayed in the required order after TaskManager failure.
   - When schema evolution is enabled on Flink, incremental.parallelism must be 
1, which is also the default. An explicitly configured value greater than one 
is rejected during job startup.
   
   ### How was this patch tested?
   
   Please Refer: MysqlCDCWithFlinkSchemaChangeIT
   
   ### Check list
   
   * [ ] If any new Jar binary package adding in your PR, please add License 
Notice according
     [New License 
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/developer/new-license.md)
   * [ ] If necessary, please update the documentation to describe the new 
feature. https://github.com/apache/seatunnel/tree/dev/docs
   * [ ] If necessary, please update `incompatible-changes.md` to describe the 
incompatibility caused by this PR.
   * [ ] If you are contributing the connector code, please check that the 
following files are updated:
     1. Update 
[plugin-mapping.properties](https://github.com/apache/seatunnel/blob/dev/plugin-mapping.properties)
 and add new connector information in it
     2. Update the pom file of 
[seatunnel-dist](https://github.com/apache/seatunnel/blob/dev/seatunnel-dist/pom.xml)
     3. Add ci label in 
[label-scope-conf](https://github.com/apache/seatunnel/blob/dev/.github/workflows/labeler/label-scope-conf.yml)
     4. Add e2e testcase in 
[seatunnel-e2e](https://github.com/apache/seatunnel/tree/dev/seatunnel-e2e/seatunnel-connector-v2-e2e/)
     5. Update connector 
[plugin_config](https://github.com/apache/seatunnel/blob/dev/config/plugin_config)
   


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