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]