DanielLeens commented on PR #12301:
URL: https://github.com/apache/seatunnel/pull/12301#issuecomment-5678515912

   Thanks @201811510411lw — this is a substantial writeup, so I went and read 
the actual diff of the referenced commit (`201811510411lw/seatunnel@c56a0f96c`) 
rather than taking the description at face value, focused on the two root-cause 
gaps from my last review.
   
   **Issue 1 (routing never reaches the real `MultiTableSink`-wrapped path)** — 
the fix shape checks out. `MultiTableSink` now implements 
`SupportSinkWriteRouting<SeaTunnelRow>` and `getWriteRouting()` aggregates each 
constituent sink's own routing policy from `sinks` (the per-table map), 
rejecting a mix of routed/unrouted tables and rejecting two source tables that 
would independently route to the same physical target. 
`SinkExecuteProcessor.createVersionSpecificDataStreamSink()` now resolves 
routing via `SupportSinkWriteRouting.resolve(sink, parallelism)` on the object 
actually passed in — which, per the call chain, is the 
post-`tryGenerateMultiTableSink` object — so this does reach the real 
production path rather than the never-constructed per-table capability the 
original PR checked. That closes the gap @goutamadwant and I both traced.
   
   **Issue 2 (schema-control rows crashing/misrouting through the data 
partitioner)** — also checks out. 
`SinkWriteRoutingPartitioner.RoutingKeySelector.getKey()` inspects 
`row.getOptions()` for `schema_change_event` *before* calling 
`routing.route(row)`, and for that case returns `schema_subtask_id` directly 
without ever invoking the routing/conversion logic that crashed on zero-field 
rows. That's exactly the "route control rows by `schema_subtask_id`, not 
through the data partitioner" fix @goutamadwant asked for, and it avoids 
`RowConverter.reconvert` seeing a zero-field row entirely rather than making 
that path merely not crash.
   
   **Issue 3 (wiring-layer regression coverage)** — 
`PaimonFlinkBucketRoutingTest`/`SinkWriteRoutingPartitionerTest` go through the 
actual `tryGenerateMultiTableSink` → `SinkExecuteProcessor` path with a real 
MiniCluster, which is the coverage gap I flagged; I haven't run these myself 
(no local execution per policy) but the shape is right — testing through the 
real wiring, not just the inner hashing logic.
   
   One scope note, consistent with what I said in my last review comment: the 
candidate is considerably larger than the three issues on this PR — 
global-commit recovery, 
`FlinkRecoverySink`/`FlinkRecoveryGlobalCommitterOperator`, new 
checkpoint-state serializers, and the Paimon writer/committer rework are all 
real additions but go beyond what's needed to close #12243's routing bug 
specifically. I'd still favor landing the routing fix (Issues 1-2, with the 
Issue 3 coverage) as the focused change here, with global-commit recovery split 
into its own PR/issue for a dedicated review, rather than merging all ~6k lines 
under this PR's original scope — but that's a call between you two on how to 
sequence it, not something I need to gate on.
   
   To be clear on where this leaves the review: no commit has landed on 
**this** PR (`48264cda6a` is unchanged, still the two dev-merge commits only), 
so I'm not doing a fresh full pass yet — this is a check of the referenced 
candidate, not a re-review of PR #12301's own head. Once the routing fix (in 
whatever scope you two agree on) is actually pushed here, I'll do a complete 
re-review against the new diff.
   


-- 
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