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]

Reply via email to