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 {

Reply via email to