This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 2e29b9b781 [flink] Fix schema inference for multi partition Kafka
topics (#9478)
2e29b9b781 is described below
commit 2e29b9b7814a7b6ae32d74f0986f9debf839ac8c
Author: Arnav Balyan <[email protected]>
AuthorDate: Wed Sep 2 13:50:11 2026 +0530
[flink] Fix schema inference for multi partition Kafka topics (#9478)
---
.../flink/action/cdc/kafka/KafkaActionUtils.java | 10 ++++-----
.../action/cdc/kafka/KafkaActionITCaseBase.java | 21 ++++++++++++++++++
.../flink/action/cdc/kafka/KafkaSchemaITCase.java | 25 ++++++++++++++++++++++
3 files changed, 50 insertions(+), 6 deletions(-)
diff --git
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionUtils.java
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionUtils.java
index da4d1a5d6e..c34c34038d 100644
---
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionUtils.java
+++
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionUtils.java
@@ -46,8 +46,6 @@ import
org.apache.kafka.common.serialization.ByteArrayDeserializer;
import java.time.Duration;
import java.util.Arrays;
-import java.util.Collection;
-import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Locale;
@@ -278,10 +276,10 @@ public class KafkaActionUtils {
+ "'topic' and 'bootstrap.servers'
config.",
topic));
}
- int firstPartition =
-
partitionInfos.stream().map(PartitionInfo::partition).sorted().findFirst().get();
- Collection<TopicPartition> topicPartitions =
- Collections.singletonList(new TopicPartition(topic,
firstPartition));
+ List<TopicPartition> topicPartitions =
+ partitionInfos.stream()
+ .map(partition -> new TopicPartition(topic,
partition.partition()))
+ .collect(Collectors.toList());
consumer.assign(topicPartitions);
consumer.seekToBeginning(topicPartitions);
diff --git
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBase.java
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBase.java
index f36dd2bf87..b2e445984e 100644
---
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBase.java
+++
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBase.java
@@ -249,6 +249,19 @@ public abstract class KafkaActionITCaseBase extends
CdcActionITCaseBase {
.forEach(r -> send(topic, r, wait));
}
+ protected void writeRecordsToKafka(
+ String topic, int partition, String resourceDirFormat, Object...
args)
+ throws Exception {
+ URL url =
+ KafkaActionITCaseBase.class
+ .getClassLoader()
+ .getResource(String.format(resourceDirFormat, args));
+ assert url != null;
+ Files.readAllLines(Paths.get(url.toURI())).stream()
+ .filter(this::isRecordLine)
+ .forEach(r -> send(topic, partition, r));
+ }
+
protected boolean isRecordLine(String line) {
try {
objectMapper.readTree(line);
@@ -269,6 +282,14 @@ public abstract class KafkaActionITCaseBase extends
CdcActionITCaseBase {
}
}
+ private void send(String topic, int partition, String record) {
+ try {
+ kafkaProducer.send(new ProducerRecord<>(topic, partition, null,
record)).get();
+ } catch (InterruptedException | ExecutionException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
/** Kafka container extension for junit5. */
private static class KafkaContainerExtension extends KafkaContainer
implements BeforeAllCallback, AfterAllCallback {
diff --git
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaSchemaITCase.java
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaSchemaITCase.java
index 84bb802bc1..be1d555cb4 100644
---
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaSchemaITCase.java
+++
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaSchemaITCase.java
@@ -67,6 +67,31 @@ public class KafkaSchemaITCase extends KafkaActionITCaseBase
{
assertThat(kafkaSchema.fields()).isEqualTo(fields);
}
+ @Test
+ @Timeout(60)
+ public void testKafkaSchemaFromNonFirstPartition() throws Exception {
+ final String topic = "test_kafka_schema_from_non_first_partition";
+ createTestTopic(topic, 2, 1);
+ writeRecordsToKafka(topic, 1,
"kafka/canal/table/schemaevolution/canal-data-1.txt");
+
+ Configuration kafkaConfig =
Configuration.fromMap(getBasicKafkaConfig());
+ kafkaConfig.setString(VALUE_FORMAT.key(), "canal-json");
+ kafkaConfig.setString(TOPIC.key(), topic);
+
+ Schema kafkaSchema =
+ MessageQueueSchemaUtils.getSchema(
+ getKafkaEarliestConsumer(
+ kafkaConfig, new
KafkaDebeziumJsonDeserializationSchema()),
+ getDataFormat(kafkaConfig),
+ TypeMapping.defaultMapping());
+
+ assertThat(kafkaSchema.fields())
+ .containsExactly(
+ new DataField(0, "pt", DataTypes.INT()),
+ new DataField(1, "_id", DataTypes.INT().notNull()),
+ new DataField(2, "v1", DataTypes.VARCHAR(10)));
+ }
+
@Test
@Timeout(60)
public void testTableOptionsChange() throws Exception {