This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-11112-70faa7ce50c62793fc2ade19df288b5219774af6 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
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\"")); + } }
