Apache SeaTunnel Kafka Sink 连接器实战指南:配置详解、消息分区与 Exactly-Once 语义
数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载Apache SeaTunnel 的 Kafka Sink 连接器用于将 SeaTunnel 任务中任意上游数据如 FakeSource、JDBC、CDC 等写入 Kafka Topic支持 Spark、Flink 与 SeaTunnel Zeta 三种引擎并提供从简单写入到 Exactly-Once 事务语义、按消息内容自定义分区、AWS MSK 与 Kerberos 安全认证等一整套生产级能力。本文以 Kafka Sink 官方文档 为主线结合 connector-kafka 模块源码 深入剖析每个配置项的真实作用与底层实现读完你即可独立编写可复制的 Kafka Sink 配置并理解其分区策略与两阶段提交2PC事务的工作原理。引擎支持与核心能力Kafka Sink 连接器基于 SeaTunnel Connector V2 API 开发可运行在以下引擎之上SparkFlinkSeaTunnel ZetaSeaTunnel 自研引擎其核心能力矩阵如下能力定义可参考 Connector V2 特性说明特性支持情况exactly-once✅ 支持cdc基于主键写入 INSERT/UPDATE/DELETE 行类型❌ 不支持关于 exactly-once官方文档特别说明默认情况下SeaTunnel 使用两阶段提交2PC来保证消息恰好一次地写入 Kafka。注意这并非意味着默认配置就是 Exactly-Once而是指当你选择EXACTLY_ONCE语义时其底层实现方式为 Kafka 事务 2PC这一点在 KafkaSemantics.java 的枚举注释中同样有明确记载AT this semantics, we will use 2pc to guarantee the message is sent to kafka exactly once.依赖与获取方式使用 Kafka 连接器需要引入对应依赖可通过install-plugin.sh脚本安装或从 Maven 中央仓库获取数据源支持版本Maven 坐标Kafka通用Universalorg.apache.seatunnel:seatunnel-connectors-v2:connector-kafkaSink 参数详解Kafka Sink 的全部配置项定义位于 Config.java下表为官方文档给出的完整参数清单参数名类型是否必填默认值说明topicString是-Sink 模式下为要写入数据的 Topic 名称bootstrap.serversString是-Kafka Broker 地址列表逗号分隔kafka.configMap否-除上述必填参数外可在此透传 Kafka Producer 客户端的任意非必填参数覆盖 Kafka 官方文档 中全部 Producer 参数semanticsString否NON可选 EXACTLY_ONCE / AT_LEAST_ONCE / NON默认 NONpartition_key_fieldsArray否-配置哪些字段作为 Kafka 消息的 KeypartitionInt否-指定分区所有消息都将发送到该分区assign_partitionsArray否-根据消息内容决定发送到哪个分区用于消息分发transaction_prefixString否-当 semantics 为 EXACTLY_ONCE 时Producer 会将所有消息写入一个 Kafka 事务Kafka 通过不同 transactionId 区分不同事务。该参数为 transactionId 的前缀务必保证不同任务使用不同前缀formatString否json数据格式默认 json可选 text、canal_json、debezium_json、ogg_json、avro 等field_delimiterString否,自定义数据格式的字段分隔符common-options-否-Sink 插件公共参数详见 Sink Common Optionstopic静态 Topic 与动态 Topic 两种格式topic支持两种写法静态 Topic直接填写 Topic 名称所有数据写入同一 Topic。动态 Topic使用上游数据中某个字段的值作为 Topic格式为${your field name}。动态 Topic 的解析逻辑位于 DefaultSeaTunnelRowSerializer.java源码使用正则\$\{(.*?)\}匹配 topic 配置若匹配成功则取出括号内的字段名从SeaTunnelRowType中查找该字段逐行将字段值作为目标 Topic 返回若配置的字段在上游 schema 中不存在会抛出Field name { ... } is not found!异常若字段值为空则抛出The column value is empty!异常。例如上游数据如下nameagedataJack16data-example1Mary23data-example2若将topic配置为${name}则第一行数据写入JackTopic第二行数据写入MaryTopic。bootstrap.serversBroker 地址列表Kafka 集群地址多个 Broker 使用英文逗号分隔例如localhost:9092,localhost:9093。在 KafkaSinkWriter.getKafkaProperties() 中该值会被设置到 Producer 的bootstrap.servers配置同时 Key/Value 序列化器均被固定为ByteArraySerializer消息先由 SeaTunnel 序列化为字节数组再交给 Kafka Producer 发送。kafka.config透传全部 Producer 参数除了上述必填参数你可以在kafka.config中以 Map 形式配置 Kafka Producer 客户端的任意参数例如acks、buffer.memory、request.timeout.ms、security.protocol、sasl.*等。源码中该 Map 会被直接平铺进 KafkaProperties对象见 KafkaSinkWriter.java因此官方 Producer 文档中的全部参数在此均可生效。本文后续的 AWS MSK 与 Kerberos 示例都依赖此机制完成安全认证配置。semantics三种投递语义与 2PC 实现semantics可选三个值默认NONEXACTLY_ONCEProducer 将全部消息写入一个 Kafka 事务在每次 checkpoint 时提交该事务从而保证每条消息恰好被写入一次AT_LEAST_ONCE在 checkpoint 时等待 Kafka Producer 缓冲区中所有未确认消息被 Kafka 确认消息不会丢失但可能重复NON不提供任何保证Kafka Broker 异常时消息可能丢失也可能产生重复。三种语义在源码中的枚举定义见 KafkaSemantics.java。当选择 EXACTLY_ONCE 时写入链路由三个类协作完成KafkaTransactionSender.java为每个 checkpoint 生成形如{transactionPrefix}-{checkpointId}的 transactionId见KafkaSinkWriter.generateTransactionId调用initTransactions()与beginTransaction()开启事务prepareCommit()时通过getProducerId()/getEpoch()保存事务现场KafkaSinkCommitter.java在 checkpoint 提交阶段通过resumeTransaction(producerId, epoch)恢复事务并执行commitTransaction()若任务失败则执行abortTransaction()回滚KafkaInternalProducer.java对 Kafka Producer 的事务管理做封装支持跨 checkpoint 恢复与中止历史事务abortTransaction(checkpointId)会从指定 checkpoint 起连续 abort 直到 epoch 归零。这正是文档中 By default, we will use 2pc to guarantee the message is sent to kafka exactly once 的实现基础。partition_key_fields以字段值作为消息 Key若希望使用上游数据的某些字段作为 Kafka 消息 Key可将字段名配置到partition_key_fields。仍以上表数据为例若将name配置为 Key 字段则消息 Key 的哈希值将决定消息进入哪个分区若不配置该参数则发送 null 消息 Key。消息 Key 的格式为 JSON例如配置name为 Key 时Key 内容形如{name:Jack}。源码层面的关键行为所选字段必须是上游 schema 中真实存在的字段KafkaSinkWriter.getPartitionKeyFields() 会逐一校验找不到字段时抛出Partition key field not found异常partition与partition_key_fields互斥同时配置会抛出Cannot select both partiton and partition_key_fields异常见 KafkaSinkWriter.getSerializer()两者的配置是二选一配置partition_key_fields时走按 Key 哈希分区配置partition时所有消息强制写入指定分区两者都不配置时所有消息发送 null Key 并随机分配分区见 KafkaSinkWriter.java。partition固定分区指定一个整数分区号后所有消息都会被发送到该分区源码中partitionExtractor(partition)直接返回固定值见 DefaultSeaTunnelRowSerializer.java。适合需要严格顺序或数据必须落盘到指定分区的场景。assign_partitions按消息内容自定义分区当 Kafka Topic 存在多个分区且希望依据消息内容做分发时可配置assign_partitions。官方示例Topic 共有 5 个分区配置assign_partitions [shoe, clothing]则包含shoe的消息发送到分区 0因为shoe在assign_partitions中下标为 0包含clothing的消息发送到分区 1其余消息通过哈希算法分配到剩余分区。该功能由MessageContentPartitioner类实现它实现了 Kafka 的org.apache.kafka.clients.producer.Partitioner接口核心逻辑见 MessageContentPartitioner.java按配置顺序遍历assign_partitions列表用message.contains(...)判断消息内容是否包含某个关键词命中第 i 个则返回分区 i全部未命中时取(message.hashCode() Integer.MAX_VALUE) % (numPartitions - assignPartitionsSize) assignPartitionsSize即哈希值落到剩余分区区间。注意该方法基于消息内容子串匹配contains且未命中消息的分区计算结果可能重叠。若你需要完全自定义的分区规则同样需要实现Partitioner接口并在kafka.config中通过partitioner.class指向你自己的实现类。transaction_prefix事务 ID 前缀当semantics EXACTLY_ONCE时Producer 的所有消息会写入一个 Kafka 事务不同事务由不同 transactionId 区分而transaction_prefix就是 transactionId 的前缀。务必为不同任务配置不同前缀避免事务 ID 冲突导致提交错乱。值得注意的是若未显式配置该参数KafkaSinkWriter 会随机生成一个形如SeaTunnel%04d如SeaTunnel1234的前缀并通过状态恢复机制在任务重启后沿用原前缀。format 与 field_delimiter消息序列化格式format默认json官方文档列出的可选值包括text、canal_json、debezium_json、ogg_json、avro。使用 json 或 text 格式时默认字段分隔符为, 可通过field_delimiter自定义默认,。若使用 canal 格式可参考 canal-json 格式说明使用 debezium 格式可参考 debezium-json 格式说明。从源码 MessageFormat.java 看实际支持的枚举比文档更丰富还包括maxwell_json、compatible_debezium_json、compatible_kafka_connect_json序列化实现位于 DefaultSeaTunnelRowSerializer.createSerializationSchema()各格式分别对应 SeaTunnel Formats 模块下的JsonSerializationSchema、TextSerializationSchema、CanalJsonSerializationSchema、OggJsonSerializationSchema、DebeziumJsonSerializationSchema、AvroSerializationSchema等。当使用compatible_debezium_json格式时动态 Topic 的取值字段会切换为上游数据中的 topic 字段见 DefaultSeaTunnelRowSerializer.java。任务示例简单示例FakeSource 写入 Kafka以下任务定义了一个 SeaTunnel 同步作业FakeSource 自动生成 16 行数据row.num16每行包含namestring与ageint两个字段最终写入test_topic该 Topic 中也将有 16 行数据。若尚未安装部署 SeaTunnel请先参考 安装部署指南 完成安装再参考 SeaTunnel 引擎快速开始 运行该作业。# Defining the runtime environment env { parallelism 1 job.mode BATCH } source { FakeSource { parallelism 1 result_table_name fake row.num 16 schema { fields { name string age int } } } } sink { kafka { topic test_topic bootstrap.servers localhost:9092 format json kafka.request.timeout.ms 60000 semantics EXACTLY_ONCE kafka.config { acks all request.timeout.ms 60000 buffer.memory 33554432 } } }几个值得注意的细节示例中同时出现了kafka.request.timeout.ms顶层参数经kafka.config之外直接透传与kafka.config内部的request.timeout.ms二者最终都会进入 Kafka Producer 配置kafka.config中的acks all表示要求所有副本确认写入是 Exactly-Once 场景下的推荐配置semantics EXACTLY_ONCE配合 checkpoint 机制通过 2PC 保证不重不丢。AWS MSK SASL/SCRAM 认证将下方${username}与${password}替换为 AWS MSK 中配置的账号密码即可sink { kafka { topic seatunnel bootstrap.servers localhost:9092 format json kafka.request.timeout.ms 60000 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}; } } }AWS MSK IAM 认证使用 IAM 认证前需要从aws-msk-iam-auth官方 release 页面下载对应版本的 jar例如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 配置如下sink { kafka { topic seatunnel bootstrap.servers localhost:9092 format json kafka.request.timeout.ms 60000 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 } } }Kerberos 认证示例对于使用 Kerberos 认证的 Kafka 集群如企业内网集群Sink 配置如下sink { Kafka { topic seatunnel bootstrap.servers 127.0.0.1:9092 format json semantics EXACTLY_ONCE kafka.config { security.protocolSASL_PLAINTEXT sasl.kerberos.service.namekafka sasl.mechanismGSSAPI java.security.krb5.conf/etc/krb5.conf sasl.jaas.configcom.sun.security.auth.module.Krb5LoginModule required \n useKeyTabtrue \n storeKeytrue \n keyTab\/path/to/xxx.keytab\ \n principal\userxxx.com\; } } }其中java.security.krb5.conf指向 Kerberos 配置文件sasl.jaas.config中的keyTab指向 keytab 密钥文件路径principal为 Kerberos 主体名均需按实际环境替换。小结与进阶阅读Kafka Sink 连接器的核心要点可归纳为四层必填参数topic、bootstrap.servers保证基本连通分发策略partition、partition_key_fields、assign_partitions决定消息落盘到哪个分区一致性语义semanticstransaction_prefix决定消息投递的可靠级别其中 EXACTLY_ONCE 由 Kafka 事务与 2PC 实现序列化与安全format、field_delimiter、kafka.config决定消息格式与认证方式。若需深入了解 Connector V2 特性定义、Canal/Debezium 等消息格式细节可继续阅读 Connector V2 特性说明、Sink Common Options 以及 canal-json、debezium-json 格式文档Kafka 连接器完整源码位于 connector-kafka 模块可对照本文逐一验证各参数的实际实现。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel Pulsar Sink 连接器实战指南参数配置、消息路由与 Exactly-Once 语义解析SeaTunnel Pulsar Sink 连接器实战指南参数配置、消息路由与 Exactly Once 语义解析 本文围绕 Apache SeaTunnel数据工程大数据批处理流处理SeaTunnel Kafka Sink 连接器完全指南配置、Exactly-Once 语义与生产实践SeaTunnel Kafka Sink 连接器完全指南配置、Exactly Once 语义与生产实践 SeaTunnel 的 Kafka Sink 连接器负数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel RocketMQ Sink Connector 深度指南配置、消息分区与 Exactly-Once 实现SeaTunnel RocketMQ Sink Connector 深度指南配置、消息分区与 Exactly Once 实现 本篇技术指南以 Apache S数据工程大数据批处理流处理上一篇终极指南如何让GitHub下载速度提升10倍以上 - Fast-GitHub浏览器插件完整教程下一篇VLC鼠标点击暂停插件一键控制视频播放的革命性工具创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考