This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new 62efde4746 [Feature][Connector-V2] Add kafka_message_value_fields to 
Kafka sink (#11112)
62efde4746 is described below

commit 62efde4746a8af32445a6a5401e903eff3d9ee16
Author: zhiliang-wu <[email protected]>
AuthorDate: Fri Sep 11 02:49:57 2026 +0000

    [Feature][Connector-V2] Add kafka_message_value_fields to Kafka sink 
(#11112)
    
    Co-authored-by: WU Zhiliang (External) <[email protected]>
---
 docs/en/connectors/sink/Kafka.md                   |  1 +
 docs/zh/connectors/sink/Kafka.md                   |  1 +
 .../seatunnel/kafka/config/KafkaSinkOptions.java   |  8 ++++
 .../serialize/DefaultSeaTunnelRowSerializer.java   | 45 +++++++++++++-------
 .../seatunnel/kafka/sink/KafkaSinkWriter.java      | 49 ++++++++++++++++++++++
 .../DefaultSeaTunnelRowSerializerTest.java         | 47 +++++++++++++++++++++
 6 files changed, 136 insertions(+), 15 deletions(-)

diff --git a/docs/en/connectors/sink/Kafka.md b/docs/en/connectors/sink/Kafka.md
index 40c8b2b0a1..08604ce68a 100644
--- a/docs/en/connectors/sink/Kafka.md
+++ b/docs/en/connectors/sink/Kafka.md
@@ -41,6 +41,7 @@ They can be downloaded via install-plugin.sh or from the 
Maven central repositor
 | semantics             | String | No       | NON     | Semantics that can be 
chosen EXACTLY_ONCE/AT_LEAST_ONCE/NON, default NON.                             
                                                                                
                                                                                
                                                                                
                                                                                
               [...]
 | partition_key_fields  | Array  | No       | -       | Configure which fields 
are used as the key of the kafka message.                                       
                                                                                
                                                                                
                                                                                
                                                                                
              [...]
 | kafka_headers_fields  | Array  | No       | -       | Configure which fields 
are used as the headers of the kafka message. The field value will be converted 
to a string and used as the header value.                                       
                                                                                
                                                                                
                                                                                
              [...]
+| kafka_message_value_fields | Array  | No       | -       | Configure which 
fields are used as the value of the kafka message. If not specified, all fields 
in the row (except those listed in `kafka_headers_fields`) will be used. Note: 
This option is not supported for `native`, `compatible_debezium_json`, and 
`compatible_kafka_connect_json` formats.                                        
                               |
 | partition             | Int    | No       | -       | We can specify the 
partition, all messages will be sent to this partition.                         
                                                                                
                                                                                
                                                                                
                                                                                
                  [...]
 | assign_partitions     | Array  | No       | -       | We can decide which 
partition to send based on the content of the message. The function of this 
parameter is to distribute information.                                         
                                                                                
                                                                                
                                                                                
                     [...]
 | transaction_prefix    | String | No       | -       | If `semantics` is 
`EXACTLY_ONCE`, the producer writes messages in Kafka transactions. Kafka 
distinguishes transactions by transaction id, so use a different prefix for 
each job.                                                                       
                                                                                
                                                        |
diff --git a/docs/zh/connectors/sink/Kafka.md b/docs/zh/connectors/sink/Kafka.md
index 1a17c2057e..393123df96 100644
--- a/docs/zh/connectors/sink/Kafka.md
+++ b/docs/zh/connectors/sink/Kafka.md
@@ -41,6 +41,7 @@ import ChangeLog from '../changelog/connector-kafka.md';
 | semantics            | String | 否    | NON  | 可以选择的语义是 
EXACTLY_ONCE/AT_LEAST_ONCE/NON,默认 NON。                                          
                                                                                
                                                                                
          |
 | partition_key_fields | Array  | 否    | -    | 配置字段用作 kafka 消息的key            
                                                                                
                                                                                
                                                                    |
 | kafka_headers_fields | Array  | 否    | -    | 配置字段用作 kafka 
消息的headers。字段值将被转换为字符串并用作 header 值                                              
                                                                                
                                                                                
     |
+| kafka_message_value_fields | Array  | 否    | -    | 配置哪些字段作为 kafka 消息的 
value。如果没有指定,则将使用行中的所有字段(除了 `kafka_headers_fields` 中的字段)。 注意:此选项不支持 `native`, 
`compatible_debezium_json` 和 `compatible_kafka_connect_json` 格式。                
                                                                      |
 | partition            | Int    | 否    | -    | 可以指定分区,所有消息都会发送到此分区            
                                                                                
                                                                                
                                                                    |
 | assign_partitions    | Array  | 否    | -    | 可以根据消息的内容决定发送哪个分区,该参数的作用是分发信息  
                                                                                
                                                                                
                                                                    |
 | transaction_prefix   | String | 否    | -    | 当 `semantics` 为 `EXACTLY_ONCE` 
时,生产者会把消息写入 Kafka 事务。Kafka 通过 transaction id 区分不同事务,因此不同作业应使用不同前缀。              
                                                                                
                                             |
diff --git 
a/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/config/KafkaSinkOptions.java
 
b/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/config/KafkaSinkOptions.java
index 05fabad51c..a23f0329fd 100644
--- 
a/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/config/KafkaSinkOptions.java
+++ 
b/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/config/KafkaSinkOptions.java
@@ -54,6 +54,14 @@ public class KafkaSinkOptions extends KafkaBaseOptions {
                             "Configure which fields are used as the headers of 
the kafka message. "
                                     + "The field value will be converted to a 
string and used as the header value.");
 
+    public static final Option<List<String>> KAFKA_MESSAGE_VALUE_FIELDS =
+            Options.key("kafka_message_value_fields")
+                    .listType()
+                    .noDefaultValue()
+                    .withDescription(
+                            "Configure which fields are used as the value of 
the kafka message. "
+                                    + "If not specified, all fields in the row 
(except headers) will be used.");
+
     public static final Option<KafkaSemantics> SEMANTICS =
             Options.key("semantics")
                     .enumType(KafkaSemantics.class)
diff --git 
a/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/serialize/DefaultSeaTunnelRowSerializer.java
 
b/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/serialize/DefaultSeaTunnelRowSerializer.java
index ec57ef4ce0..397437feec 100644
--- 
a/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/serialize/DefaultSeaTunnelRowSerializer.java
+++ 
b/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/serialize/DefaultSeaTunnelRowSerializer.java
@@ -129,13 +129,14 @@ public class DefaultSeaTunnelRowSerializer implements 
SeaTunnelRowSerializer {
             MessageFormat format,
             String delimiter,
             ReadonlyConfig pluginConfig) {
-        return create(topic, partition, null, rowType, format, delimiter, 
pluginConfig);
+        return create(topic, partition, null, null, rowType, format, 
delimiter, pluginConfig);
     }
 
     public static DefaultSeaTunnelRowSerializer create(
             String topic,
             Integer partition,
             List<String> headerFields,
+            List<String> messageValueFields,
             SeaTunnelRowType rowType,
             MessageFormat format,
             String delimiter,
@@ -145,7 +146,8 @@ public class DefaultSeaTunnelRowSerializer implements 
SeaTunnelRowSerializer {
                 partitionExtractor(partition),
                 timestampExtractor(),
                 keyExtractor(null, rowType, format, delimiter, pluginConfig),
-                valueExtractor(headerFields, rowType, format, delimiter, 
pluginConfig),
+                valueExtractor(
+                        headerFields, messageValueFields, rowType, format, 
delimiter, pluginConfig),
                 headersExtractor(headerFields, rowType));
     }
 
@@ -156,13 +158,14 @@ public class DefaultSeaTunnelRowSerializer implements 
SeaTunnelRowSerializer {
             MessageFormat format,
             String delimiter,
             ReadonlyConfig pluginConfig) {
-        return create(topic, keyFields, null, rowType, format, delimiter, 
pluginConfig);
+        return create(topic, keyFields, null, null, rowType, format, 
delimiter, pluginConfig);
     }
 
     public static DefaultSeaTunnelRowSerializer create(
             String topic,
             List<String> keyFields,
             List<String> headerFields,
+            List<String> messageValueFields,
             SeaTunnelRowType rowType,
             MessageFormat format,
             String delimiter,
@@ -172,7 +175,8 @@ public class DefaultSeaTunnelRowSerializer implements 
SeaTunnelRowSerializer {
                 partitionExtractor(null),
                 timestampExtractor(),
                 keyExtractor(keyFields, rowType, format, delimiter, 
pluginConfig),
-                valueExtractor(headerFields, rowType, format, delimiter, 
pluginConfig),
+                valueExtractor(
+                        headerFields, messageValueFields, rowType, format, 
delimiter, pluginConfig),
                 headersExtractor(headerFields, rowType));
     }
 
@@ -310,18 +314,21 @@ public class DefaultSeaTunnelRowSerializer implements 
SeaTunnelRowSerializer {
 
     private static Function<SeaTunnelRow, byte[]> valueExtractor(
             List<String> headerFields,
+            List<String> messageValueFields,
             SeaTunnelRowType rowType,
             MessageFormat format,
             String delimiter,
             ReadonlyConfig pluginConfig) {
-        if (headerFields == null || headerFields.isEmpty()) {
+        if ((headerFields == null || headerFields.isEmpty())
+                && (messageValueFields == null || 
messageValueFields.isEmpty())) {
             return valueExtractor(rowType, format, delimiter, pluginConfig);
         }
 
-        // Create a new row type excluding header fields
-        SeaTunnelRowType valueRowType = createValueRowType(headerFields, 
rowType);
+        // Create a new row type excluding header fields or retaining only 
message value fields
+        SeaTunnelRowType valueRowType =
+                createValueRowType(headerFields, messageValueFields, rowType);
         Function<SeaTunnelRow, SeaTunnelRow> valueRowExtractor =
-                createValueRowExtractor(valueRowType, headerFields, rowType);
+                createValueRowExtractor(valueRowType, rowType);
         SerializationSchema serializationSchema =
                 createSerializationSchema(valueRowType, format, delimiter, 
false, pluginConfig);
         return row -> 
serializationSchema.serialize(valueRowExtractor.apply(row));
@@ -345,16 +352,24 @@ public class DefaultSeaTunnelRowSerializer implements 
SeaTunnelRowSerializer {
     }
 
     private static SeaTunnelRowType createValueRowType(
-            List<String> headerFieldNames, SeaTunnelRowType rowType) {
-        // Create a row type excluding header fields
+            List<String> headerFieldNames,
+            List<String> messageValueFields,
+            SeaTunnelRowType rowType) {
         List<String> valueFieldNames = new java.util.ArrayList<>();
         List<SeaTunnelDataType> valueFieldTypes = new java.util.ArrayList<>();
 
-        for (int i = 0; i < rowType.getTotalFields(); i++) {
-            String fieldName = rowType.getFieldName(i);
-            if (!headerFieldNames.contains(fieldName)) {
+        if (messageValueFields != null && !messageValueFields.isEmpty()) {
+            for (String fieldName : messageValueFields) {
                 valueFieldNames.add(fieldName);
-                valueFieldTypes.add(rowType.getFieldType(i));
+                
valueFieldTypes.add(rowType.getFieldType(rowType.indexOf(fieldName)));
+            }
+        } else {
+            for (int i = 0; i < rowType.getTotalFields(); i++) {
+                String fieldName = rowType.getFieldName(i);
+                if (headerFieldNames == null || 
!headerFieldNames.contains(fieldName)) {
+                    valueFieldNames.add(fieldName);
+                    valueFieldTypes.add(rowType.getFieldType(i));
+                }
             }
         }
 
@@ -383,7 +398,7 @@ public class DefaultSeaTunnelRowSerializer implements 
SeaTunnelRowSerializer {
     }
 
     private static Function<SeaTunnelRow, SeaTunnelRow> 
createValueRowExtractor(
-            SeaTunnelRowType valueType, List<String> headerFieldNames, 
SeaTunnelRowType rowType) {
+            SeaTunnelRowType valueType, SeaTunnelRowType rowType) {
         int[] valueIndex = new int[valueType.getTotalFields()];
         for (int i = 0; i < valueType.getTotalFields(); i++) {
             valueIndex[i] = rowType.indexOf(valueType.getFieldName(i));
diff --git 
a/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaSinkWriter.java
 
b/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaSinkWriter.java
index 9fcbc91098..cacadc9d19 100644
--- 
a/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaSinkWriter.java
+++ 
b/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaSinkWriter.java
@@ -61,6 +61,7 @@ import static 
org.apache.seatunnel.connectors.seatunnel.kafka.config.KafkaSinkOp
 import static 
org.apache.seatunnel.connectors.seatunnel.kafka.config.KafkaSinkOptions.FORMAT;
 import static 
org.apache.seatunnel.connectors.seatunnel.kafka.config.KafkaSinkOptions.KAFKA_CONFIG;
 import static 
org.apache.seatunnel.connectors.seatunnel.kafka.config.KafkaSinkOptions.KAFKA_HEADERS_FIELDS;
+import static 
org.apache.seatunnel.connectors.seatunnel.kafka.config.KafkaSinkOptions.KAFKA_MESSAGE_VALUE_FIELDS;
 import static 
org.apache.seatunnel.connectors.seatunnel.kafka.config.KafkaSinkOptions.PARTITION;
 import static 
org.apache.seatunnel.connectors.seatunnel.kafka.config.KafkaSinkOptions.PARTITION_KEY_FIELDS;
 import static 
org.apache.seatunnel.connectors.seatunnel.kafka.config.KafkaSinkOptions.SEMANTICS;
@@ -182,6 +183,19 @@ public class KafkaSinkWriter implements 
SinkWriter<SeaTunnelRow, KafkaCommitInfo
             ReadonlyConfig pluginConfig, SeaTunnelRowType seaTunnelRowType) {
         MessageFormat messageFormat = pluginConfig.get(FORMAT);
         String topic = pluginConfig.get(TOPIC);
+
+        if (pluginConfig.get(KAFKA_MESSAGE_VALUE_FIELDS) != null) {
+            if (MessageFormat.NATIVE.equals(messageFormat)
+                    || 
MessageFormat.COMPATIBLE_DEBEZIUM_JSON.equals(messageFormat)
+                    || 
MessageFormat.COMPATIBLE_KAFKA_CONNECT_JSON.equals(messageFormat)) {
+                throw new KafkaConnectorException(
+                        CommonErrorCode.OPERATION_NOT_SUPPORTED,
+                        String.format(
+                                "kafka_message_value_fields is not supported 
for %s format",
+                                messageFormat));
+            }
+        }
+
         if (MessageFormat.NATIVE.equals(messageFormat)) {
             // Validate that kafka_headers_fields is not configured for NATIVE 
format
             if (pluginConfig.get(KAFKA_HEADERS_FIELDS) != null) {
@@ -207,6 +221,7 @@ public class KafkaSinkWriter implements 
SinkWriter<SeaTunnelRow, KafkaCommitInfo
         // Validate that partition_key_fields and kafka_headers_fields don't 
overlap
         List<String> partitionKeyFields = getPartitionKeyFields(pluginConfig, 
seaTunnelRowType);
         List<String> headerFields = getHeaderFields(pluginConfig, 
seaTunnelRowType);
+        List<String> messageValueFields = getMessageValueFields(pluginConfig, 
seaTunnelRowType);
         if (!partitionKeyFields.isEmpty() && !headerFields.isEmpty()) {
             for (String headerField : headerFields) {
                 if (partitionKeyFields.contains(headerField)) {
@@ -218,12 +233,25 @@ public class KafkaSinkWriter implements 
SinkWriter<SeaTunnelRow, KafkaCommitInfo
                 }
             }
         }
+        // Validate that kafka_message_value_fields and kafka_headers_fields 
don't overlap
+        if (!messageValueFields.isEmpty() && !headerFields.isEmpty()) {
+            for (String headerField : headerFields) {
+                if (messageValueFields.contains(headerField)) {
+                    throw new KafkaConnectorException(
+                            CommonErrorCode.ILLEGAL_ARGUMENT,
+                            String.format(
+                                    "Field '%s' cannot be in both 
kafka_message_value_fields and kafka_headers_fields",
+                                    headerField));
+                }
+            }
+        }
 
         if (pluginConfig.get(PARTITION_KEY_FIELDS) != null) {
             return DefaultSeaTunnelRowSerializer.create(
                     topic,
                     partitionKeyFields,
                     headerFields,
+                    messageValueFields,
                     seaTunnelRowType,
                     messageFormat,
                     delimiter,
@@ -234,6 +262,7 @@ public class KafkaSinkWriter implements 
SinkWriter<SeaTunnelRow, KafkaCommitInfo
                     topic,
                     pluginConfig.get(PARTITION),
                     headerFields,
+                    messageValueFields,
                     seaTunnelRowType,
                     messageFormat,
                     delimiter,
@@ -244,6 +273,7 @@ public class KafkaSinkWriter implements 
SinkWriter<SeaTunnelRow, KafkaCommitInfo
                 topic,
                 Collections.<String>emptyList(),
                 headerFields,
+                messageValueFields,
                 seaTunnelRowType,
                 messageFormat,
                 delimiter,
@@ -308,6 +338,25 @@ public class KafkaSinkWriter implements 
SinkWriter<SeaTunnelRow, KafkaCommitInfo
         return Collections.emptyList();
     }
 
+    private List<String> getMessageValueFields(
+            ReadonlyConfig pluginConfig, SeaTunnelRowType seaTunnelRowType) {
+        if (pluginConfig.get(KAFKA_MESSAGE_VALUE_FIELDS) != null) {
+            List<String> messageValueFields = 
pluginConfig.get(KAFKA_MESSAGE_VALUE_FIELDS);
+            List<String> rowTypeFieldNames = 
Arrays.asList(seaTunnelRowType.getFieldNames());
+            for (String messageValueField : messageValueFields) {
+                if (!rowTypeFieldNames.contains(messageValueField)) {
+                    throw new KafkaConnectorException(
+                            CommonErrorCode.ILLEGAL_ARGUMENT,
+                            String.format(
+                                    "Message value field not found: %s, 
rowType: %s",
+                                    messageValueField, rowTypeFieldNames));
+                }
+            }
+            return messageValueFields;
+        }
+        return Collections.emptyList();
+    }
+
     private void checkNativeSeaTunnelType(SeaTunnelRowType seaTunnelRowType) {
         SeaTunnelRowType exceptRowType = 
nativeTableSchema().toPhysicalRowDataType();
         for (int i = 0; i < exceptRowType.getFieldTypes().length; i++) {
diff --git 
a/seatunnel-connectors-v2/connector-kafka/src/test/java/org/apache/seatunnel/connectors/seatunnel/kafka/serialize/DefaultSeaTunnelRowSerializerTest.java
 
b/seatunnel-connectors-v2/connector-kafka/src/test/java/org/apache/seatunnel/connectors/seatunnel/kafka/serialize/DefaultSeaTunnelRowSerializerTest.java
index ec3575139b..374a18d9e7 100644
--- 
a/seatunnel-connectors-v2/connector-kafka/src/test/java/org/apache/seatunnel/connectors/seatunnel/kafka/serialize/DefaultSeaTunnelRowSerializerTest.java
+++ 
b/seatunnel-connectors-v2/connector-kafka/src/test/java/org/apache/seatunnel/connectors/seatunnel/kafka/serialize/DefaultSeaTunnelRowSerializerTest.java
@@ -95,6 +95,7 @@ public class DefaultSeaTunnelRowSerializerTest {
                         topic,
                         Arrays.asList("id"),
                         Arrays.asList("source", "traceId"),
+                        null,
                         rowType,
                         format,
                         delimiter,
@@ -138,6 +139,7 @@ public class DefaultSeaTunnelRowSerializerTest {
                         topic,
                         Arrays.asList("id"),
                         Arrays.asList("source", "traceId"),
+                        null,
                         rowType,
                         format,
                         delimiter,
@@ -251,6 +253,7 @@ public class DefaultSeaTunnelRowSerializerTest {
                         topic,
                         Arrays.asList("id"),
                         Arrays.asList("source", "traceId"),
+                        null,
                         rowType,
                         format,
                         delimiter,
@@ -307,6 +310,7 @@ public class DefaultSeaTunnelRowSerializerTest {
                         topic,
                         Arrays.asList("id"),
                         Arrays.asList("source", "traceId"),
+                        null,
                         rowType,
                         format,
                         delimiter,
@@ -335,4 +339,47 @@ public class DefaultSeaTunnelRowSerializerTest {
         Assertions.assertFalse(valueString.contains("\"source\""));
         Assertions.assertFalse(valueString.contains("\"traceId\""));
     }
+
+    @Test
+    public void testMessageValueFields() {
+        String topic = "test_topic";
+        SeaTunnelRowType rowType =
+                new SeaTunnelRowType(
+                        new String[] {"id", "name", "source", "traceId"},
+                        new 
org.apache.seatunnel.api.table.type.SeaTunnelDataType[] {
+                            BasicType.INT_TYPE,
+                            BasicType.STRING_TYPE,
+                            BasicType.STRING_TYPE,
+                            BasicType.STRING_TYPE
+                        });
+        MessageFormat format = MessageFormat.JSON;
+        String delimiter = ",";
+        Map<String, Object> configMap = new HashMap<>();
+        ReadonlyConfig pluginConfig = ReadonlyConfig.fromMap(configMap);
+
+        // Test with message value fields
+        DefaultSeaTunnelRowSerializer serializer =
+                DefaultSeaTunnelRowSerializer.create(
+                        topic,
+                        Arrays.asList("id"), // partition_key_fields
+                        null, // header_fields
+                        Arrays.asList("name", "source"), // 
message_value_fields
+                        rowType,
+                        format,
+                        delimiter,
+                        pluginConfig);
+
+        SeaTunnelRow row = new SeaTunnelRow(new Object[] {1, "test", "web", 
"trace-123"});
+        ProducerRecord<byte[], byte[]> record = serializer.serializeRow(row);
+
+        Assertions.assertEquals("test_topic", record.topic());
+
+        String valueString = new String(record.value(), 
StandardCharsets.UTF_8);
+        // The value should only contain the requested fields
+        Assertions.assertTrue(valueString.contains("\"name\""));
+        Assertions.assertTrue(valueString.contains("\"source\""));
+        // The id and traceId should be excluded
+        Assertions.assertFalse(valueString.contains("\"id\""));
+        Assertions.assertFalse(valueString.contains("\"traceId\""));
+    }
 }

Reply via email to