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\""));
+ }
}