goutamadwant commented on code in PR #12301:
URL: https://github.com/apache/seatunnel/pull/12301#discussion_r4000362920


##########
seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/main/java/org/apache/seatunnel/core/starter/flink/execution/SinkExecuteProcessor.java:
##########
@@ -64,9 +69,24 @@ protected DataStreamSink<SeaTunnelRow> 
createVersionSpecificDataStreamSink(
                             .name("BroadcastSchemaHandler")
                             .setParallelism(parallelism);
         }
+        if (sink instanceof SupportSinkDataPartition) {

Review Comment:
   This check runs after `tryGenerateMultiTableSink`. Paimon implements 
`SupportMultiTableSink`, so even a one-table Paimon job is wrapped in 
`MultiTableSink`; that wrapper does not implement `SupportSinkDataPartition`. 
The normal Flink Paimon path therefore never enters this branch, and the 
reported fixed-bucket corruption remains. The same applies to Paimon sinks 
created from `tables_configs`.
   
   I reran the #12264 reproduction on this PR head. Passing `PaimonSink` 
directly made all eight cases pass. Changing only the test to obtain the sink 
through `tryGenerateMultiTableSink`, as the production processor does, made all 
six two-writer cases fail again: 100 rows remained, the deleted key was 
present, and the update stayed at `before`, both without and after a completed 
checkpoint.
   
   Please propagate per-table partition ownership through `MultiTableSink`, 
account for `multi_table_sink_replica` when determining physical writer 
ownership, and add the regression through the normal sink-creation path.



##########
seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/main/java/org/apache/seatunnel/core/starter/flink/execution/SinkExecuteProcessor.java:
##########
@@ -64,9 +69,24 @@ protected DataStreamSink<SeaTunnelRow> 
createVersionSpecificDataStreamSink(
                             .name("BroadcastSchemaHandler")
                             .setParallelism(parallelism);
         }
+        if (sink instanceof SupportSinkDataPartition) {
+            Optional<SinkDataPartitioner<SeaTunnelRow>> partitioner =
+                    ((SupportSinkDataPartition<SeaTunnelRow>) sink)
+                            .getSinkDataPartitioner(parallelism);
+            if (partitioner.isPresent()) {
+                ds = partitionBySinkDataPartitioner(ds, partitioner.get());
+            }
+        }
         return ds.sinkTo(
                         SinkV1Adapter.wrap(
                                 new FlinkSink<>(sink, 
stream.getCatalogTables(), parallelism)))
                 .name(String.format("%s-Sink", sink.getPluginName()));
     }
+
+    private DataStream<SeaTunnelRow> partitionBySinkDataPartitioner(
+            DataStream<SeaTunnelRow> stream, SinkDataPartitioner<SeaTunnelRow> 
partitioner) {
+        return stream.partitionCustom(

Review Comment:
   `partitionCustom` also receives the zero-field schema-control rows emitted 
by `BroadcastSchemaSinkOperator`, before `FlinkSinkWriter` can consume them. 
The key selector invokes the Paimon partitioner on that row. I reproduced this 
using the emitted `new SeaTunnelRow(0)` shape: it throws 
`ArrayIndexOutOfBoundsException` in `RowConverter.reconvert` at 
`SeaTunnelRow.getField(0)`.
   
   Even if conversion were skipped, repartitioning these control rows can send 
multiple schema events to one writer while another sink subtask never applies 
or acknowledges the change. Please preserve one schema-control event per 
downstream subtask—for example, route these rows by the existing 
`schema_subtask_id` rather than invoking the sink data partitioner—and add a 
fixed-bucket schema-evolution regression.



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