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]