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]
