Kafka Sink 接收器完整实战指南

发布时间:2026/9/19 1:25:39
Kafka Sink 接收器完整实战指南 Kafka Sink 接收器完整实战指南【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel一、为什么需要 Kafka SinkSeaTunnel 与 Kafka 的数据出口SeaTunnel 的 Kafka Sink 接收器Sink将上游 Source 产生的 SeaTunnel Rows 内容写入 Kafka topic是海量数据集成管道中最常用的数据出口之一。它原生支持 Spark、Flink 与 SeaTunnel Zeta 三种引擎并内置了基于 Kafka 事务的两阶段提交2PC精确一次语义可满足从简单日志投递到金融级不重不丢的流式管道等多种场景。读完本篇你将掌握Kafka Sink 的全部配置参数与默认值、EXACTLY_ONCE/AT_LEAST_ONCE/NON三种语义的底层实现原理、基于字段值路由 topic 与分区含MessageContentPartitioner的实战用法、Kafka Headers 与消息 Value 的字段裁剪、json/text/canal_json/debezium_json/protobuf/NATIVE等格式的选择以及 AWS MSK SASL/SCRAM、AWS MSK IAM、Kerberos 三种安全认证配置。二、支持引擎与核心特性支持引擎Spark / Flink / SeaTunnel Zeta精确一次Exactly-Once✅ 已支持CDC❌ 不支持定时刷新❌ 不支持默认情况下连接器使用 2PC 保证消息只发送一次到 Kafka。该能力由连接器源码中的KafkaSink实现它实现了SeaTunnelSinkSeaTunnelRow, KafkaSinkState, KafkaCommitInfo, KafkaAggregatedCommitInfo接口通过createWriter/restoreWriter创建KafkaSinkWriter并通过createCommitter创建KafkaSinkCommitter来完成事务提交源码见 KafkaSink.java。三、支持的数据源信息与依赖使用 Kafka 连接器需要以下依赖可通过install-plugin.sh安装或从 Maven 中央存储库下载数据源支持版本MavenKafka通用org.apache.seatunnel:connector-kafka安装后Kafka 相关 JAR 会位于$SEATUNNEL_HOME/plugin/kafka/lib目录下见下文 AWS MSK IAM 认证示例中aws-msk-iam-authJAR 的放置位置。四、接收器选项Sink Options全表名称类型是否必需默认值描述topicString是-当表用作接收器时topic 名称是要写入数据的 topicbootstrap.serversString是-Kafka brokers使用逗号分隔kafka.configMap否-除上述 Kafka Producer 客户端必须指定的参数外用户还可以为 Producer 客户端指定多个非强制参数涵盖 Kafka 官方文档中指定的所有生产者参数semanticsString否NON可选语义EXACTLY_ONCE / AT_LEAST_ONCE / NON默认 NONpartition_key_fieldsArray否-配置字段用作 Kafka 消息的 keykafka_headers_fieldsArray否-配置字段用作 Kafka 消息的 headers字段值将被转换为字符串并用作 header 值kafka_message_value_fieldsArray否-配置哪些字段作为 Kafka 消息的 value未指定时使用行中所有字段除kafka_headers_fields中的字段注意此选项不支持native、compatible_debezium_json和compatible_kafka_connect_json格式partitionInt否-可以指定分区所有消息都会发送到此分区assign_partitionsArray否-可以根据消息内容决定发送哪个分区作用是分发消息transaction_prefixString否-当semantics为EXACTLY_ONCE时生产者会把消息写入 Kafka 事务。Kafka 通过 transaction id 区分不同事务因此不同作业应使用不同前缀formatString否json数据格式默认 json。可选 text、canal_json、debezium_json、compatible_debezium_json、ogg_json、maxwell_json、avro、protobuf 和 native。若使用 json 或 text 格式默认字段分隔符是,自定义分隔符请添加field_delimiter选项。canal 格式参考 canal-jsondebezium 格式参考 debezium-jsonfield_delimiterString否,自定义数据格式的字段分隔符common-options否-Sink 插件常用参数详见 Sink 常用选项protobuf_message_nameString否-format 配置为 protobuf 时生效取 Message 名称protobuf_schemaString否-format 配置为 protobuf 时生效取 Schema 名称上述参数在源码中均有明确定义KafkaSinkOptions.java、KafkaBaseOptions.java。其中semantics被定义为枚举类型KafkaSemantics默认值为NONformat被定义为枚举类型MessageFormat默认值为JSON。五、参数详解5.1 Topic 格式目前支持两种写法直接填写 topic 名称例如topic test_topic。使用上游数据中的字段值作为 topic格式为${your field name}其中 topic 是上游数据某一列的值。例如上游数据如下nameagedataJack16data-example1Mary23data-example2如果设置topic ${name}则第一行发送到 Jack topic第二行发送到 Mary topic。5.2 语义Semantics在源码中KafkaSemantics枚举定义了三种语义KafkaSemantics.javaEXACTLY_ONCE生产者将在 Kafka 事务中写入所有消息这些消息在检查点checkpoint上提交给 Kafka。该模式能保证数据精确写入 Kafka 一次即使任务失败重试也不会出现数据重复和丢失。AT_LEAST_ONCE生产者将等待 Kafka 缓冲区中所有未完成的消息在检查点上被 Kafka 生产者确认。该模式保证数据至少写入 Kafka 一次即使任务失败。NON不提供任何保证。如果 Kafka 代理出现问题消息可能丢失也可能重复任务失败重试可能产生数据丢失或重复。使用EXACTLY_ONCE时需要开启 checkpoint并确保每个运行中的作业使用唯一的transaction_prefix。多个作业复用同一个事务前缀可能导致 Kafka 事务冲突。从源码看实现在KafkaSinkWriter的构造函数中如果semantics为EXACTLY_ONCE则创建KafkaTransactionSender基于 Kafka 事务的 2PC 发送器否则创建KafkaNoTransactionSenderKafkaSinkWriter.java。具体行为KafkaTransactionSender为每个 checkpoint 生成形如{transaction_prefix}-{checkpointId}的事务 IDgenerateTransactionId方法在snapshotState时开启下一批事务在prepareCommit时 flush 并检查异步发送失败将KafkaCommitInfo含 transactionId、producerId、epoch、txnStarted 标记交给KafkaSinkCommitter在 Broker 上真正commitTransactionKafkaTransactionSender.java、KafkaSinkCommitter.java。KafkaNoTransactionSender直接sendsnapshotState时仅flushbeginTransaction/abortTransaction/prepareCommit均为 no-opKafkaNoTransactionSender.java。另外源码中transaction_prefix未配置时会随机生成SeaTunnel%04d格式的前缀范围 0000-9999因此生产环境建议显式配置以保证作业间的唯一性。5.3 分区关键字段partition_key_fields如果你想使用上游数据中的字段值作为 Kafka 消息的 key可以将这些字段名指定给此属性。上游数据如下所示nameagedataJack16data-example1Mary23data-example2如果将name设置为 key那么name列的哈希值将决定消息发送到哪个分区。消息 key 的格式为 json如果设置name为 key例如{name:Jack}。所选的字段必须是上游数据中已存在的字段否则作业启动时会抛出Partition key field not found异常源码校验逻辑见KafkaSinkWriter.getPartitionKeyFields。注意partition与partition_key_fields二者只能配置其一。若同时配置源码会抛出Cannot select both partition and partition_key_fields异常。另外partition_key_fields与kafka_headers_fields不能有重叠字段。5.4 Kafka Headers 字段kafka_headers_fields如果你想使用上游数据中的字段值作为 Kafka 消息的 headers可以将这些字段名指定给此属性。上游数据如下所示nameagedatasourcetraceIdJack16data-example1webtrace-123Mary23data-example2mobiletrace-456如果将source和traceId设置为 Kafka headers 字段这些字段值将作为 headers 添加到 Kafka 消息中。例如第一行将具有 headerssourceweb和traceIdtrace-123。字段值将被转换为字符串并用作 header 值所选的字段必须是上游数据中已存在的字段。注意两点配置为 Kafka headers 的字段不会包含在消息的 valuepayload中只会存在于 Kafka 消息的 headers 中。源码中同时校验了kafka_headers_fields与kafka_message_value_fields不可重叠Field %s cannot be in both ...。format native时不支持kafka_headers_fields配置会抛出kafka_headers_fields is not supported with NATIVE format异常。5.5 分区分配assign_partitions假设 topic 总有五个分区配置中的assign_partitions字段设置为assign_partitions [shoe, clothing]在这种情况下包含 shoe 的消息将被发送到第零个分区因为 shoe 在assign_partitions中被标记为下标 0包含 clothing 的消息将被发送到第一个分区。对于其他消息将使用哈希算法将它们均匀地分配到剩余的分区中。该功能通过MessageContentPartitioner类实现它实现了org.apache.kafka.clients.producer.Partitioner接口MessageContentPartitioner.java。从源码看其partition方法按顺序遍历assign_partitions列表对消息 value 做contains匹配命中第 i 个即返回分区 i未命中的消息通过HashUtils.bucketIndex(message.hashCode(), numPartitions - assignPartitionsSize) assignPartitionsSize落到剩余分区。如果需要自定义分区策略可以仿照该类实现Partitioner接口并通过kafka.config的partitioner.class参数指定。六、任务示例6.1 简单示例FakeSource 写入 Kafka此示例展示了如何定义一个 SeaTunnel 同步任务通过 FakeSource 自动产生数据并发送到 Kafka Sink。FakeSource 会生成总共 16 行数据row.num16每一行包含name字符串和age整型两个字段。最终这些数据被发送到 test_topic该 topic 也将包含 16 行数据。如果还未安装和部署 SeaTunnel请参照 安装 SeaTunnel 指南部署完成后可按 快速开始使用 SeaTunnel 引擎 运行任务。# Defining the runtime environment env { parallelism 1 job.mode BATCH } source { FakeSource { parallelism 1 plugin_output fake row.num 16 schema { fields { name string age int } } } } sink { kafka { topic test_topic bootstrap.servers localhost:9092 format json semantics EXACTLY_ONCE kafka.config { acks all request.timeout.ms 60000 buffer.memory 33554432 } } }源码中getKafkaProperties会将kafka.config中的键值对原样放入Properties并强制设置key.serializer与value.serializer为ByteArraySerializer消息统一按字节数组传输具体编码交给format决定。6.2 使用 Kafka Headers本示例展示如何使用kafka_headers_fields将上游字段设置为 Kafka 消息头同时通过partition_key_fields指定消息 keyenv { parallelism 1 job.mode BATCH } source { FakeSource { parallelism 1 plugin_output fake row.num 16 schema { fields { name string age int source string traceId string } } } } sink { Kafka { topic test_topic bootstrap.servers localhost:9092 format json partition_key_fields [name] kafka_headers_fields [source, traceId] semantics EXACTLY_ONCE kafka.config { acks all request.timeout.ms 60000 buffer.memory 33554432 } } }6.3 AWS MSK SASL/SCRAM 认证将以下${username}和${password}替换为 AWS MSK 中的配置值sink { kafka { topic seatunnel bootstrap.servers localhost:9092 format json semantics EXACTLY_ONCE kafka.config { security.protocolSASL_SSL sasl.mechanismSCRAM-SHA-512 sasl.jaas.configorg.apache.kafka.common.security.scram.ScramLoginModule required \nusername${username}\npassword${password}; } } }6.4 AWS MSK IAM 认证从 AWS MSK IAM Auth 官方 release 页面下载aws-msk-iam-auth-1.1.5.jar放入$SEATUNNEL_HOME/plugin/kafka/lib目录。请确保 IAM 策略具有kafka-cluster:Connect权限如下配置Effect: Allow, Action: [ kafka-cluster:Connect, kafka-cluster:AlterCluster, kafka-cluster:DescribeCluster ],接收器配置sink { kafka { topic seatunnel bootstrap.servers localhost:9092 format json semantics EXACTLY_ONCE kafka.config { security.protocolSASL_SSL sasl.mechanismAWS_MSK_IAM sasl.jaas.configsoftware.amazon.msk.auth.iam.IAMLoginModule required; sasl.client.callback.handler.classsoftware.amazon.msk.auth.iam.IAMClientCallbackHandler } } }6.5 Kerberos 认证示例请在启动 SeaTunnel 之前设置 JVM 参数java.security.krb5.conf或更新/etc/krb5.conf中的默认krb5.conf。接收器配置示例sink { Kafka { topic seatunnel bootstrap.servers localhost:9092 format json semantics EXACTLY_ONCE kafka.config { security.protocol SASL_PLAINTEXT sasl.kerberos.service.name kafka sasl.mechanism GSSAPI sasl.jaas.config com.sun.security.auth.module.Krb5LoginModule required \n useKeyTabtrue \n storeKeytrue \n keyTab\/path/to/xxx.keytab\ \n principal\userxxx.com\; } } }6.6 Protobuf 配置将format设置为protobuf并配置protobuf_schema与protobuf_message_name参数对应源码 KafkaBaseOptions.java 中的PROTOBUF_SCHEMA、PROTOBUF_MESSAGE_NAMEsink { kafka { topic test_protobuf_topic_fake_source bootstrap.servers kafkaCluster:9092 format protobuf kafka.config { acks all request.timeout.ms 60000 buffer.memory 33554432 } protobuf_message_name Person protobuf_schema syntax proto3; package org.apache.seatunnel.format.protobuf; option java_outer_classname ProtobufE2E; message Person { int32 c_int32 1; int64 c_int64 2; float c_float 3; double c_double 4; bool c_bool 5; string c_string 6; bytes c_bytes 7; message Address { string street 1; string city 2; string state 3; string zip 4; } Address address 8; mapstring, float attributes 9; repeated string phone_numbers 10; } } }6.7 NATIVE 格式写入 Kafka 原生消息如果需要写入 Kafka 原生的信息消息自带 key/partition/timestamp/headers可以参考下面的配置。sink { kafka { topic test_topic_native_sink bootstrap.servers kafkaCluster:9092 format NATIVE } }输入参数要求如下key/value需要byte[]类型{ headers: { header1: header1, header2: header2 }, key: dGVzdF9ieXRlc19kYXRh, partition: 3, timestamp: 1672531200000, timestampType: CREATE_TIME, value: dGVzdF9ieXRlc19kYXRh }从源码看NATIVE 格式要求上游 Schema 严格匹配固定的原生表结构headersmapstring,string、keybyte[]、partitionint、timestamplong、valuebyte[]任何字段缺失或类型不符都会在作业启动时报错KafkaSinkWriter.checkNativeSeaTunnelType。6.8 流式 EXACTLY_ONCE 与 Checkpoint 协同长时间运行的流式作业若不能容忍丢失或重复需要开启 checkpoint 并把semantics设置为EXACTLY_ONCE。SeaTunnel 会把 Kafka 事务与 checkpoint 协调保证每个 in-flight 批次和对应的消费 offset 原子提交。env { parallelism 2 job.mode STREAMING checkpoint.interval 10000 } source { Kafka { topic orders bootstrap.servers localhost:9092 consumer.group orders_consumer start_mode group_offsets format json schema { fields { order_id bigint user_id bigint amount double } } } } sink { Kafka { topic orders_sink bootstrap.servers localhost:9092 format json semantics EXACTLY_ONCE transaction_prefix orders_pipeline kafka.config { transaction.timeout.ms 900000 } partition_key_fields [order_id] } }注意每个作业都必须使用唯一的transaction_prefix。Kafka 通过 transactional id 区分事务跨作业复用相同前缀会导致事务冲突。建议transaction.timeout.ms配置为略大于 checkpoint 间隔避免事务在 checkpoint 提交前超时。6.9 NATIVE 格式下的 Header 转发当 source 端使用format NATIVE、sink 端同样使用format NATIVE时Kafka 的 headers、key、partition、timestamp 等字段会原样回写到下游 topic。这种用法适合在不改写数据布局的前提下把已是 Kafka 编码的记录转发到另一个 topic。source { Kafka { topic topic_native_source bootstrap.servers localhost:9092 format NATIVE consumer.group native_forwarder } } sink { Kafka { topic topic_native_sink bootstrap.servers localhost:9092 format NATIVE } }注意上游记录使用format NATIVE时key和value是byte[]。这种情况下请谨慎配置kafka_headers_fields因为 headers 已经编码在行内。七、常见问题FAQ7.1 Kafka Sink 会自动创建 topic 吗SeaTunnel Kafka Sink 本身不会主动创建 Kafka topic只是向配置的topic写入数据。topic 是否自动创建取决于 Kafka Broker 的auto.create.topics.enable配置。生产环境中建议提前手动创建 topic以便自行控制分区数、副本数、保留策略和 ACL。不要依赖自动创建因为 Broker 可能已将auto.create.topics.enable设为false。7.2 不配置partition_key_fields会怎样若未设置partition_key_fieldsSeaTunnel 将以null作为 Kafka 消息 key 发送记录源码中对应Collections.StringemptyList()分支Kafka 会使用默认的轮询策略将记录分散到各分区。这适合做负载均衡但不适合需要相同业务 key 的记录落入同一分区以保证顺序的场景。如有顺序要求请配置partition_key_fields。7.3 如何实现精确一次exactly-once写入将semantics设为EXACTLY_ONCE以启用精确一次语义并配置transaction_prefix为每个任务提供唯一的 Kafka 事务 ID 前缀。SeaTunnel 会将 Kafka 事务与 checkpoint 协调来实现 exactly-oncesink { kafka { topic output-topic bootstrap.servers localhost:9092 semantics EXACTLY_ONCE transaction_prefix SeaTunnelJob kafka.config { transaction.timeout.ms 900000 } } }确保 Kafka Broker 开启了事务支持且transaction.timeout.ms与 checkpoint 间隔相匹配。在EXACTLY_ONCE语义下发送失败会让 checkpoint 失败而不是静默丢弃数据。此时可能出现两种错误错误码名称含义处理建议KAFKA-08TRANSACTION_NOT_STARTED事务中已有数据但 Kafka 始终未在 Broker 端完成该事务的注册。检查 Broker 是否可用以及transaction.timeout.ms是否小于 checkpoint 间隔。KAFKA-09PRODUCE_DATA_FAILED事务中的某条数据异步发送失败。查看异常 cause可重试的异常通常在 checkpoint 重试后恢复其他异常需要排查 Broker 端问题。从源码看这两种错误分别对应 KafkaConnectorErrorCode.java 中的TRANSACTION_NOT_STARTED与PRODUCE_DATA_FAILED前者在prepareCommit阶段发现recordNumInTransaction 0但isTxnStarted()仍为 false 时主动抛出后者由KafkaTransactionSender通过AtomicReferenceException asyncSendException记录事务内的首个异步发送失败并在写入路径/提交路径上抛出。两种错误都会中止当前事务受影响的数据会从上一个已完成的 checkpoint 重新发送不会丢失。7.4 如何配置 SASL/Kerberos 认证sink { kafka { topic secure-topic bootstrap.servers broker:9092 kafka.security.protocol SASL_PLAINTEXT kafka.sasl.mechanism GSSAPI kafka.sasl.kerberos.service.name kafka kafka.sasl.jaas.config com.sun.security.auth.module.Krb5LoginModule required useKeyTabtrue keyTab/etc/kafka/kafka.keytab principaluserREALM.COM; } }7.5 Kafka Sink 支持哪些消息格式支持json、text、canal_json、debezium_json、ogg_json、avro、protobuf和NATIVE源码枚举MessageFormat还包含compatible_debezium_json、compatible_kafka_connect_json、maxwell_json见 MessageFormat.java。当上游数据已经是带 headers、key 和 value 字节字段的 Kafka 原生格式时使用NATIVE格式。八、变更日志Kafka 连接器的历史变更记录见 connector-kafka 变更日志。九、延伸阅读Sink 常用选项所有 Sink 插件的公共参数。canal-json 格式 与 debezium-json 格式CDC 场景下 JSON 消息格式详解。SeaTunnel 安装部署 与 SeaTunnel 引擎快速开始本地环境准备与任务提交。Kafka 连接器测试用例KafkaTransactionSenderTest.java、MessageContentPartitionerTest.java可了解分区分配与事务语义的验证细节。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考