SeaTunnel 实战:用 Kafka Source + Iceberg Sink 搭建流式事件入湖链路
SeaTunnel 实战用 Kafka Source Iceberg Sink 搭建流式事件入湖链路【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本篇技术指南以 Apache SeaTunnel 中「Kafka 到 Iceberg」的经典链路为核心完整讲解如何把 Kafka 里的流式事件实时落到 Iceberg 表中供后续 Spark、Trino 等引擎分析查询。你将掌握连接器插件安装、最小 HOCON 配置、本地模式启动流任务、结果验证以及常见排错手段并理解 upsert、分区、schema 演进等关键选项在 Iceberg Sink 底层的作用机制。链路概览为什么选择 Kafka → Iceberg当业务希望把 Kafka 中的流式事件订单、点击、日志等落成一份可被分析引擎直接查询的湖仓表时Iceberg 是常见选择。Kafka 承担实时消息缓冲Iceberg 则提供 ACID 事务、时间旅行、schema 演进等表能力。在 SeaTunnel 中这条链路由两个连接器协作完成Kafka Source订阅 topic、按声明的 schema 反序列化消息Iceberg Sink负责自动建表、分区写入、upsert 提交与 schema 演进。从仓库源码看Iceberg Sink 的能力覆盖 CDC 写入、自动建表、表结构变更与多表写入见 Iceberg.md 的「描述」与「主要特性」因此 Kafka 中的 JSON 事件经过一次简单的 source/sink 配置即可平稳入湖。前置条件1. 先跑通第一个任务建议先完成 跑第一个任务用 FakeSource Console 验证本地部署、配置解析与执行引擎均正常避免把环境问题混入链路调试。2. 安装 Kafka 与 Iceberg 连接器插件SeaTunnel 从 2.2.0-beta 起二进制包不再默认附带连接器依赖。部署说明见 部署 下载连接器插件将 config/plugin_config 收敛为仅包含本链路所需插件--seatunnel-connectors-- connector-kafka connector-iceberg --end--然后执行安装命令并确认插件已就位cd ${SEATUNNEL_HOME} sh bin/install-plugin.sh ls connectors | rg connector-(kafka|iceberg)install-plugin.sh 在 Linux/macOS 上通过 HTTPS 直接下载 JAR 及校验文件需要curl、mktemp以及sha512sum/sha1sum/shasum/openssl之一Windows 的 install-plugin.cmd 仍走内置 Maven Wrapper具体见 deployment.md。3. Flink/Spark 引擎的特殊依赖如果你使用 Flink 或 Spark 引擎运行该链路需要补齐 Iceberg 在对应环境中的依赖例如hive-exec和libfb303。其原因可从 Iceberg.md 的「数据库依赖」一节确认Iceberg 连接器 pom 中hive-exec的依赖范围为providedFlink 用户需将hive-exec-xxx.jar、libfb303-xxx.jar放入FLINK_HOME/libSpark 若已集成 Hadoop 则无需额外添加。部分版本的hive-exec不内嵌libfb303需手动补充。使用 SeaTunnel 内置的 Zeta 引擎则无此负担。4. 准备 Iceberg warehouse 目录本教程使用本地 Hadoop catalogwarehouse 对应file:///tmp/seatunnel/iceberg/warehouse-demo需先创建一个当前进程可写的空目录mkdir -p /tmp/seatunnel/iceberg/warehouse-demo注意warehouse 路径必须对运行任务的引擎进程可写这是最常见的失败点之一详见文末「常见坑」。5. 准备 Kafka topic 与测试数据创建 topic 并写入两条 JSON 订单消息kafka-topics.sh \ --create \ --if-not-exists \ --topic orders \ --bootstrap-server kafka:9092 \ --partitions 1 \ --replication-factor 1 kafka-console-producer.sh --topic orders --bootstrap-server kafka:9092 EOF {id:1001,customer_id:2001,total_amount:19.99,event_date:2026-06-12} {id:1002,customer_id:2002,total_amount:29.99,event_date:2026-06-12} EOF最小配置解析将下面的配置保存为config/kafka-to-iceberg.confenv { parallelism 1 job.mode STREAMING checkpoint.interval 5000 } source { Kafka { plugin_output orders_kafka topic orders bootstrap.servers kafka:9092 consumer.group seatunnel-orders start_mode earliest format json schema { fields { id bigint customer_id bigint total_amount decimal(16, 2) event_date string } } } } sink { Iceberg { plugin_input orders_kafka catalog_name seatunnel_demo namespace lakehouse table orders iceberg.catalog.config { type hadoop warehouse file:///tmp/seatunnel/iceberg/warehouse-demo } iceberg.table.primary-keys id iceberg.table.partition-keys event_date iceberg.table.upsert-mode-enabled true iceberg.table.schema-evolution-enabled true case_sensitive true } }env 块流式任务的三要素parallelism 1本示例单并发足够生产可按 topic 分区数放大job.mode STREAMING声明为流任务任务将持续运行checkpoint.interval 5000每 5 秒触发一次 checkpoint。它既是容错恢复的基础也直接决定 Kafka 消费位点的提交时机——SeaTunnel 在 checkpoint 完成时才向 Kafka 提交 offset详见 Kafka.md 的「消费组 offset 是如何提交的」。流作业不开启 checkpoint 会导致重启后的一致性行为变弱这是常见坑之一。source 块Kafka 的关键选项选项示例值说明plugin_outputorders_kafka命名本插件输出的数据流供下游plugin_input引用topicorders订阅的主题逗号分隔可订阅多个主题bootstrap.serverskafka:9092必填Kafka brokers 列表consumer.groupseatunnel-orders消费者组 ID默认SeaTunnel-Consumer-Groupstart_modeearliest初始消费模式可选earliest/latest/group_offsets/specific_offsets/timestampformatjson消息格式默认即json还支持text、canal_json、debezium_json、ogg_json、avro、protobuf、nativeschema字段类型声明定义反序列化后的行结构关于start_mode的选择earliest从头重放全量数据group_offsets从消费组已提交位点恢复适合中断重启latest只消费启动后新消息。注意一个细节从 checkpoint 或 savepoint 恢复时Kafka Source 会优先使用 checkpoint 中保存的 split offsetstart_mode与消费组位点只在首次启动或新发现分区时生效见 Kafka.md 的「源选项」说明。total_amount声明为decimal(16, 2)而非double是为了在湖表中保留精确的小数语义——根据 Iceberg.md 的「数据类型映射」SeaTunnel 的 DECIMAL 会直接映射为 Iceberg 的 DECIMAL而浮点类型无法做到这一点。sink 块Iceberg 的核心选项选项示例值说明plugin_inputorders_kafka接入上游数据流catalog_nameseatunnel_democatalog 名称默认defaultnamespacelakehouseIceberg 数据库命名空间默认defaulttableorders目标表名不配置时使用上游表名iceberg.catalog.configtypehadoop warehouse必填初始化 Catalog 的属性可参考 Iceberg 的 CatalogPropertiesiceberg.table.primary-keysid主键列逗号分隔多个iceberg.table.partition-keysevent_date建表时的分区字段逗号分隔多个也可用 Iceberg transform 如days(ts)iceberg.table.upsert-mode-enabledtrue启用 upsert 模式默认falseiceberg.table.schema-evolution-enabledtrue允许同步过程中支持 schema 变更默认falsecase_sensitivetrue列名匹配是否区分大小写从源码看这些选项在 IcebergSinkOptions.java 中逐一声明iceberg.table.primary-keys无默认值且当iceberg.table.upsert-mode-enabled为 true 时必须显式提供主键列表——因为 upsert 模式不再自动继承 source 表主键源码注释中明确引用了该行为变更。iceberg.table.write-props可透传write.format.default、write.target-file-size-bytes等 Iceberg 表属性优先级最高iceberg.table.commit-branch可将提交写到指定分支schema_save_mode默认CREATE_SCHEMA_WHEN_NOT_EXISTdata_save_mode默认APPEND_DATA需要先删后写或自定义清理 SQL 时可分别调整。底层写入路径源码视角在 IcebergSink.java 中IcebergSinkConfig负责解析插件配置主键、分区键、upsert 开关等sink 的提交链路由IcebergAggregatedCommitter与IcebergFilesCommitter组成写入器则依据是否分区、是否启用 upsert 在PartitionedAppendWriter/PartitionedDeltaWriter/UnpartitionedDeltaWriter之间选择见 sink/writer 目录。也就是说开启 upsert 后数据以主键为基准执行 delta 合并写schema-evolution-enabled则在同步过程中把上游新增字段同步为 Iceberg 表的 schema 变更。运行任务在${SEATUNNEL_HOME}下用本地模式启动cd ${SEATUNNEL_HOME} ./bin/seatunnel.sh --config ./config/kafka-to-iceberg.conf -m local这是一条流式任务Kafka 消息被消费、写入 Iceberg 并提交的整个过程中任务必须保持运行。停止任务即停止消费这与批任务「跑完即退」的行为不同请务必留意。验证结果1. 检查 warehouse 目录Iceberg 表会在首次写入时自动建表元数据metadata目录与数据文件data目录都会落在 warehouse 下ls /tmp/seatunnel/iceberg/warehouse-demo/lakehouse/orders2. 用兼容引擎查询使用 Spark、Trino 或其他 Iceberg 兼容引擎验证数据。以 Spark SQL 为例spark-sql \ --conf spark.sql.catalog.seatunnel_demoorg.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.seatunnel_demo.typehadoop \ --conf spark.sql.catalog.seatunnel_demo.warehousefile:///tmp/seatunnel/iceberg/warehouse-demo \ -e SELECT COUNT(*) FROM seatunnel_demo.lakehouse.orders如果表可以正常查询且行数与写入 Kafka 的消息数量一致说明链路已打通。按照示例中的两条订单消息最终行数应为2。由于示例开启了 upsert 模式且主键为id若向orderstopic 写入重复id的消息行数不会线性增长而是按主键去重合并——这是验证 upsert 生效的直观手段。常见坑与排查JSON 结构与 schema 不一致Kafka 消息里的字段与 source 声明的schema必须对齐包括字段名、类型与精度。缺字段会反序列化失败类型不匹配如字符串混入 bigint 字段会在运行期抛错。必要时可借助format_error_handle_way skip跳过脏数据默认fail。流作业未开启 checkpointcheckpoint.interval缺失会使任务在故障重启后难以恢复到一致的消费位点Kafka 消费位点的提交也随之失去锚点端到端一致性大打折扣。catalog 类型正确但 warehouse 不可写Hadoop catalog 直接读写 warehouse 路径路径必须对当前引擎进程可写且磁盘要有足够空间。本地file://与 HDFS 路径hdfs://your_cluster/...均要求对应权限生产环境可参考 Iceberg.md 中的 Hadoop catalog 与 Hive catalog 示例。开启 upsert 但消息没有稳定主键upsert 模式要求每条消息都携带稳定的主键字段iceberg.table.primary-keys否则同一行的多次更新无法正确合并甚至产生数据丢失或重复。如果 Kafka 消息天然无主键应关闭 upsert默认即关闭改走纯 append 追加模式。相关文档Kafka Source 连接器完整的源选项表、start_mode语义、分区动态发现、SASL/Kerberos 认证示例Iceberg Sink 连接器全部 Sink 选项、Hive/Hadoop Catalog 示例、分支提交、Kerberos 认证与多表写入示例部署与插件安装install-plugin.sh的下载机制与镜像配置跑第一个任务本地基础链路的验证起点【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考