DanielLeens commented on issue #11752: URL: https://github.com/apache/seatunnel/issues/11752#issuecomment-5247867114
Thanks for reducing this to the parallel-sink schema-evolution path. I checked the current Flink translation side, and this does look like a real correctness issue. In `BroadcastSchemaSinkOperator`, the subtask that receives the broadcasted schema-change signal forwards a schema row downstream to its sink side and relies on the writer path to ACK after the ALTER TABLE work completes. What is still missing from the current design is a clear guarantee that every relevant parallel sink writer refreshes its local schema before post-change records can reach it. So the important part of the fix is not only "broadcast the DDL event", but also: 1. ensure all relevant parallel writers update their local schema state; 2. preserve ordering so post-change records cannot overtake that refresh; 3. avoid duplicate or unsafe external schema application while doing the above. An E2E regression with sink parallelism greater than 1 is especially important here, because this kind of bug can pass accidentally when the evolved records happen to route to the one writer that already refreshed. I see the issue is already assigned to you, which matches the PR willingness flag. Please continue from that path. -- 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]
