201811510411lw commented on PR #12301: URL: https://github.com/apache/seatunnel/pull/12301#issuecomment-5675466116
Thanks @DanielLeens for the follow-up review. I’ve published the complete candidate so the implementation and tests are available for inspection: - [Candidate commit: c56a0f96c](https://github.com/201811510411lw/seatunnel/commit/c56a0f96c72c0e5024cbfef461f391f01c9fa99c) - [Branch: fix/issue-12243-paimon-routing-recovery](https://github.com/201811510411lw/seatunnel/tree/fix/issue-12243-paimon-routing-recovery) It is one commit on Apache `dev` base `46e58fc16f8523c82edd40b93e8f8d62663737a2`. Here is how it relates to the review findings. **1. Routing through the actual MultiTableSink path** `MultiTableSink` forwards the per-table routing capability. Paimon determines ownership from the physical partition and fixed bucket, and the writer validates ownership before writing. The integration fixture obtains the sink through `tryGenerateMultiTableSink` and invokes the Flink sink processor; it explicitly asserts that even the single-table sink is wrapped. For `multi_table_sink_replica`, this candidate requires `1` for fixed-bucket sinks and rejects unsupported values before writing. It also rejects independent source-table writers targeting the same physical Paimon table, rather than allowing conflicting owners. **2. Schema-control routing** `SinkWriteRoutingPartitioner` routes zero-field control rows by `schema_subtask_id`, bypassing data routing and row conversion. Its tests verify one destination per sink subtask, no interaction with the data-routing policy, and rejection of missing or invalid destinations. To be precise about coverage: this is a partitioner-level regression, not a full online schema-evolution integration test. The complete candidate currently rejects online structural changes for fixed-bucket writers; a restore event confirming the same source schema is accepted. I am not claiming that this closes the request for end-to-end online DDL coverage. **3. Tests available in the commit** - `PaimonFlinkBucketRoutingTest`: real Flink MiniCluster/Paimon tests through the normal factory/processor path; one and two writers; INSERT/UPDATE/DELETE before and after completed checkpoints; multiple tables, partitions and buckets; full key/value readback; non-draining savepoint restoration from two writers to two or one, followed by further changes and checkpoints. - `PaimonFlinkPendingRecoveryTest`: failure after checkpoint completion but before local commit messages reach the global sink; recovery of all pending writer contributions, followed by further writes and checkpoints. - `SinkWriteRoutingPartitionerTest`: the schema-control cases described above and data-writer ownership validation. - `MultiTableWriteRoutingTest`, `PaimonCheckpointCommitTest`, and the Flink state/operator tests: routing constraints, commit completeness/order, serialization, and incompatible or incomplete recovery state. Validation used JDK 8 and Flink 1.18.1 through the Flink 1.15 adapter. The selected affected-module suite passed 610 tests across 70 classes, with zero failures/errors/skips. Scoped formatting checks and the 32-module build passed; 34 packaged implementation classes were checked against the compiled classes. This is local verification, not external MySQL/object-store acceptance or a claim that all SeaTunnel E2E tests passed. **Scope and integration** This is a complete candidate for reference, not a drop-in commit for this PR. It uses `SupportSinkWriteRouting` plus a separate `SupportSinkGlobalCommitRecovery` capability. I’m happy to align routing with `SupportSinkDataPartition` and prepare a focused follow-up rather than introduce two competing interfaces. The recovery portion ensures historical global commits are complete before writers resume. I agree that routing and its regression coverage should be the priority here; recovery can be considered separately after checking that the resulting routing change remains correct across recovery. The candidate also has explicit compatibility limits: fixed and non-fixed bucket tables cannot share this multi-table sink; old-protocol checkpoints/savepoints are not migrated; only same-parallelism recovery and validated reduction to one writer are supported. The separate Flink 1.13/1.20 adapters reject this fixed-bucket path. These restrictions are documented in the Paimon and incompatible-changes docs in both languages, and would need consideration when agreeing on the final scope. -- 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]
