CloverDew commented on PR #11960: URL: https://github.com/apache/seatunnel/pull/11960#issuecomment-5420684593
> > > I think this overlaps with some of the work already being done in #11912, especially around replacing `LocalSchemaCoordinator` and checkpoint-based schema coordination. Could you clarify how you see the relationship between this issue and #11912? Is this intended as an alternative design, or as additional work on top of it? > > > > > > It seems you focused on rewriting the SinkWriter layer. I replaced the entire Flink Coordinator, primarily addressing fault recovery issues by replacing the coordinator with a data plane protocol. I consider this additional work, and you can focus on the SinkWriter side. Your implementation essentially still utilizes the coordinator, while I completely removed it, using Flink's native mechanisms for coordination. In this branch, I didn't make any additional rewrites to the SinkWriter. > > Thanks for the clarification. Just to clarify one point about #11912: the current implementation is not limited to the SinkWriter layer, and it no longer uses LocalSchemaCoordinator. SchemaOperator also removes requestSchemaChange() and relies on Flink checkpoint completion to coordinate schema-change dispatch and the release of buffered rows. > > I agree that #11960 adds stronger recovery/rescaling mechanisms, such as explicit producer/sequence IDs, atomic protocol state, and replay handling. > > However, I think there is still an architectural overlap between the two PRs. #11912 keeps parallel sink writers for the same table and separates one-time external DDL application from writer-local schema refresh, while #11960 appears to partition data by table and let one table owner handle the schema change. So these seem less like independent SinkWriter/coordinator changes and more like two different coordination models for part of the same problem. > > Would it make sense to first agree on which parallel-writer model we want to keep, and then treat the additional recovery/rescaling guarantees in #11960 as follow-up hardening on top of that model? I agree that the two PRs should now be treated as alternative coordination models rather than independent changes. One point I would clarify is that I don't think the recovery/rescaling guarantees in #11960 are entirely orthogonal hardening that can simply be layered on top of either parallel-writer model. In #11960, deterministic table ownership, producer/sequence IDs, dependency tracking on data rows, pending controls/rows, and checkpointed protocol state are designed together as one consistency model. The ownership model is what removes the need for cross-writer coordination for a table, while the explicit protocol state makes that model recoverable and replayable after failure or rescaling. So I think the architectural choice is slightly broader than only "which parallel-writer model to keep". It is more like: #11912: keep multiple writers for the same table, execute external DDL once in a dedicated operator, broadcast the resulting schema, and refresh writer-local state. #11960: deterministically assign each table to an owner on the sink side and model schema evolution as an in-band, checkpointed protocol with explicit dependency/replay state. I agree that we should first reach consensus on which coordination model we want for Flink. But I think we should choose one of the architectures, once that is decided, we can then see which pieces from the other PR are still useful or can be reused. -- 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]
