Skip to main content

概述

Apache Iceberg Sink Connector 将 Kafka Topic 中的业务记录持续写入 Iceberg 表,连接实时事件流与数据湖中的历史查询、离线分析和审计链路。记录的值映射为表字段,写入的数据文件通过表快照提交后对查询可见;Kafka 记录的键不会自动成为表列。 一个 Connector 可以写入单张表、多张固定表,或按记录字段选择目标表。不配置路由字段时,同一条记录会写入所有固定目标表;配置路由后,可将不同类别的业务事件分开存储。它也支持自动建表和受支持的 Schema 演进,适合持续追加的订单事件、行为日志和变更历史。内置写入路径采用追加语义,不会按业务主键自动合并当前状态。 Connector 使用 Kafka 事务、控制 Topic 和 Iceberg 快照检查点协调写入与恢复,以提供有前提的 exactly-once 恢复语义。该语义要求事务协议和 Connect 运行时兼容、真实数据消费组对齐、Catalog 支持原子快照提交,并保留恢复所需的控制事件、消费组位点、数据文件和快照历史;它不等于业务主键去重,也不提供多表原子提交。

前置条件

  • 准备可用的 Iceberg Catalog,以及所选 Catalog 和 FileIO 所需的客户端依赖、认证方式和存储权限。
  • 默认关闭自动建表,须预建需要接收记录的目标表,且记录字段与目标表 Schema 兼容;静态路由实际选中的目标缺表会使 Task 失败,动态路由缺表则跳过写入。
  • 启用自动建表时,Catalog 须允许建表,目标命名空间须可用或允许创建,首条记录须能推导出非空表 Schema。
  • 准备控制 Topic,或允许 Broker 自动创建;内部客户端须具有控制 Topic 读写、事务 ID 和相关消费组操作所需权限。
  • 内部 Kafka 客户端须连接到数据消费组所在的 Kafka 集群,以便在同一 Kafka 事务中提交控制事件和输入位点。
  • Kafka 须支持包含消费组元数据的事务协议,且 Connect 运行时须兼容 Connector 获取底层消费者的方式。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka、Iceberg REST Catalog 和一张目标表,确认网络连通及 Catalog、存储、控制 Topic 的访问权限。集群与 Connector 的管理操作见管理 Connector。以下配置将不带 Schema 信封的 JSON 对象追加到已有表。
替换输入 Topic、完整表名、REST Catalog 地址和 Kafka 地址。该示例使用默认控制 Topic control-iceberg,须提前保证其可用;REST Catalog 及其存储访问配置按实际部署补充到 iceberg.catalog.*,内部 Kafka 客户端需要的安全属性补充到 iceberg.kafka.*,不要在配置中暴露长期凭证。输入值须为 JSON 对象且与已有表字段兼容;示例不适用于直接读取 Debezium 信封来维护当前状态。 默认每 5 分钟发起一次表提交。消费位点可能先于表快照推进,因此消费 Lag 降低并不意味着查询已经可见全部数据。

配置

以下配置适用于追加写入。已注册配置的默认值与实际运行回退分别说明;动态属性未单独注册统一默认值或重要级别,其校验由读取逻辑或对应组件完成。单表属性中的 <table-name> 是完整目标表名,包含命名空间。 前缀属性的默认值和重要级别中的 N/A 表示未单独定义,不表示所传递的每项属性都没有默认值;具体默认值及约束由对应组件决定。单表属性保留其实际读取回退。

输入与任务

connector.class

选择 Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 必填:是
  • 有效值 / 注意事项:使用 org.apache.iceberg.connect.IcebergSinkConnector

tasks.max

设置最大任务数。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:至少为 1;有效并行度受输入分区分配限制,不是每张表分配一个 Task,也不保证跨分区全局顺序。

topics

指定输入 Kafka Topic 列表。
  • 类型list
  • 默认值:空列表 []
  • 重要级别:高
  • 有效值 / 注意事项:逗号分隔,与非空 topics.regex 二选一;它不是目标表列表或控制 Topic。

topics.regex

按正则表达式订阅输入 Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:合法 Java 正则表达式,与非空 topics 二选一。

key.converter

将 Kafka 消息键转换为 Connect 数据。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:未设置时继承 Worker 转换器;显式类须实现 Converter。消息键不会自动写入表字段。

value.converter

将 Kafka 消息值转换为 Connect 数据。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:未设置时继承 Worker 转换器,没有统一的 JSON 默认值。转换及 SMT 处理后的非空顶层值须为 StructMap;转换器子属性须与实际输入编码匹配。

consumer.override.*

覆盖 Connect Task 的数据消费者属性。 此属性集合没有统一默认值,具体属性的类型和默认值由 Kafka Consumer 决定。
  • 类型Kafka Consumer 透传属性
  • 默认值:N/A
  • 重要级别:N/A
  • 有效值 / 注意事项:须通过 Worker 的客户端覆盖策略;优先级为 Connect 基值、Worker consumer.*、Connector consumer.override.*。覆盖 group.id 时须同步 iceberg.connect.group-id;不会配置内部控制客户端。auto.offset.reset 仅影响没有有效位点的情况,不能作为无损恢复手段。

Catalog 与存储

iceberg.catalog

设置 Catalog 实例名称,而非实现类型。
  • 类型string
  • 默认值iceberg
  • 重要级别:中
  • 有效值 / 注意事项:即使使用默认名称,也须提供非空 iceberg.catalog.* 属性集合。

iceberg.catalog.*

将去掉前缀后的属性交给 Iceberg Catalog,包括实现选择、服务地址、认证和 FileIO 配置。 此属性集合没有统一默认值,具体默认值由 Catalog 和 FileIO 决定。
  • 类型string 透传属性
  • 默认值:N/A
  • 重要级别:N/A
  • 必填:是
  • 有效值 / 注意事项typecatalog-impl 互斥。type 不区分大小写,可选择 hivehadooprestgluenessiejdbc;实际可用性取决于已安装的实现和客户端依赖。两者都未设置时底层回退到 Hive,但这不是 Connector 的默认类型声明。uriwarehouseio-impl 等是否必填由选定实现决定。

iceberg.hadoop-conf-dir

指定 Hadoop XML 配置目录。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:目录须能被 Worker 访问,读取其中的 core-site.xmlhdfs-site.xmlhive-site.xml;需要可用的 Hadoop 类。iceberg.hadoop.* 在 XML 加载后覆盖对应属性。

iceberg.hadoop.*

将去掉前缀后的属性设置到 Hadoop Configuration。 此属性集合没有统一默认值,具体默认值由 Hadoop 决定。
  • 类型string 透传属性
  • 默认值:N/A
  • 重要级别:N/A
  • 有效值 / 注意事项:需要 Hadoop 类;它不是 Kafka 客户端或 Catalog 属性前缀。

目标表与路由

iceberg.tables

设置静态目标表列表。
  • 类型list
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:静态模式必填,使用逗号分隔的完整表名;与 iceberg.tables.dynamic-enabled=true 互斥。空列表没有写入目标。不配置路由字段时,每条记录写入所有静态目标表。

iceberg.tables.dynamic-enabled

启用按记录字段选择目标表的动态路由。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:为 true 时须移除 iceberg.tables 并设置 iceberg.tables.route-field。路由值转为小写后作为完整表名,不使用单表路由正则;空值会跳过,非法表标识符可能导致失败。

iceberg.tables.route-field

指定记录值中的路由字段。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:动态模式必填;支持点分嵌套 StructMap 字段路径。静态模式下,该字段的值与每张表的 route-regex 独立匹配。字段缺失或为空时不写入任何表。

iceberg.table.<table-name>.route-regex

设置单张静态表的路由匹配条件。
  • 类型string
  • 默认值:未注册;读取回退为 null
  • 重要级别:N/A
  • 有效值 / 注意事项:静态模式设置路由字段时,须为接收数据的表提供正则;使用 Java 正则全字符串匹配,重叠表达式会将记录写入多表。没有匹配表达式的表不接收记录;动态模式忽略此项。

自动建表与 Schema

iceberg.tables.auto-create-enabled

允许创建缺失的目标表。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:根据首条记录的 Schema 或值推导表结构,需要建表权限。空对象无法推导;关闭时,静态路由实际选中的目标缺表会抛出 NoSuchTableException 并使 Task 进入 FAILED 状态,仅动态路由对该缺表异常跳过写入。动态缺表跳过不适用于非法表标识符、权限错误或其他 Catalog 异常。

iceberg.tables.evolve-schema-enabled

允许受支持的表 Schema 更新。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:可增加能推导类型的新字段;带 Schema 的 Struct 还支持 intlongfloatdouble 的提升和必填字段放宽为可选。无 Schema 数据不具有同样的已有字段提升路径;关闭时额外字段会被忽略。它不支持任意重命名、删列或类型迁移。

iceberg.tables.schema-force-optional

将创建或演进的字段设为可选。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:也作用于创建或演进的容器值类型;不会整体改写已有表所有字段的可选性。默认尊重输入 Schema 的可选性。

iceberg.tables.schema-case-insensitive

控制表字段查找是否忽略大小写。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:不重命名字段;表存在名称映射时使用映射查找路径,此开关不会无条件覆盖映射。

iceberg.tables.default-partition-by

设置自动建表的默认分区规则。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:仅作用于新建表,不改变已有表分区。支持字段本身以及 yearmonthdayhourbuckettruncate 变换;多个规则用逗号分隔,变换括号内的逗号保留,buckettruncate 的参数顺序为列名、宽度。解析或构建异常时会记录错误并退回无分区建表,不能只凭配置判断分区已生效。

iceberg.table.<table-name>.partition-by

覆盖单表自动建表分区规则。
  • 类型string
  • 默认值:未注册;读取回退为 iceberg.tables.default-partition-by
  • 重要级别:N/A
  • 有效值 / 注意事项:仅用于自动建表;显式空字符串使新表无分区。语法和异常回退与全局分区规则相同。

iceberg.tables.auto-create-props.*

为自动创建的表提供 Iceberg 表属性。 此属性集合没有统一默认值,具体默认值由 Iceberg 表和 Catalog 决定。
  • 类型Iceberg 表属性透传
  • 默认值:N/A
  • 重要级别:N/A
  • 有效值 / 注意事项:去掉前缀后传给建表操作;不会更新已有表属性。

iceberg.tables.write-props.*

覆盖写入器初始化使用的 Iceberg 表属性,例如文件格式和目标文件大小。 此属性集合没有统一默认值,具体默认值由 Iceberg 表属性决定。
  • 类型Iceberg 表属性透传
  • 默认值:N/A
  • 重要级别:N/A
  • 有效值 / 注意事项:覆盖写入器本地属性,不持久化修改表属性;目标文件大小不保证每次提交都达到该大小。

标识列与提交分支

iceberg.tables.default-id-columns

设置默认标识列列表。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:以逗号分隔并去掉两侧空白,列须存在;未设置或为空时不覆盖表 Schema 的标识字段。此项不是主键唯一性约束,也不会启用 upsert 或删除已有行。

iceberg.table.<table-name>.id-columns

覆盖单表标识列列表。
  • 类型string
  • 默认值:未注册;读取回退为 iceberg.tables.default-id-columns
  • 重要级别:N/A
  • 有效值 / 注意事项:显式空字符串抑制全局列表覆盖并保留表 Schema 标识字段;其他约束与全局项相同,不改变追加语义。

iceberg.tables.default-commit-branch

设置默认表快照提交分支。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:未设置时提交到 main。分支不必预先存在,底层快照提交允许创建新分支;已存在的同名引用必须是分支而不是 Tag。单表配置优先。

iceberg.table.<table-name>.commit-branch

覆盖单表提交分支。
  • 类型string
  • 默认值:未注册;读取回退为 iceberg.tables.default-commit-branch
  • 重要级别:N/A
  • 有效值 / 注意事项:最终未设置时使用 main;新分支和 Tag 的约束与全局项相同。写入其他分支的数据不会自动出现在 main 查询中。

控制通道与提交

iceberg.control.topic

指定传递写入文件引用和提交协调事件的控制 Topic。
  • 类型string
  • 默认值control-iceberg
  • 重要级别:中
  • 有效值 / 注意事项:不是输入数据 Topic;Connector 不主动发出创建 Topic 请求。保留周期须覆盖尚未完成提交及恢复所需的事件,不能只按正常提交周期设定。

iceberg.control.group-id-prefix

设置内部控制消费组前缀。
  • 类型string
  • 默认值cg-control-
  • 重要级别:低
  • 有效值 / 注意事项:默认值包含末尾连字符;它与 Connect 数据消费组不同,恢复依赖控制消费组位点的保留。

iceberg.control.commit.interval-ms

设置发起表提交的周期,单位毫秒。
  • 类型int
  • 默认值300000
  • 重要级别:中
  • 有效值 / 注意事项:使用正值;周期影响查询可见性与文件关闭频率,但不是端到端延迟上限。配置定义没有范围或与超时之间的大小关系校验。

iceberg.control.commit.timeout-ms

设置协调器等待写入任务响应的超时,单位毫秒。
  • 类型int
  • 默认值30000
  • 重要级别:中
  • 有效值 / 注意事项:使用正值;超时可提交已收到的部分响应,不意味着 Task 必然失败,也不是全部 Catalog 提交操作的总耗时上限。

iceberg.control.commit.threads

设置协调器并行提交目标表的线程数。
  • 类型int
  • 默认值Runtime.getRuntime().availableProcessors() * 2
  • 重要级别:中
  • 有效值 / 注意事项:默认表达式在配置定义初始化时求值,使用 Worker JVM 可见处理器数量的两倍。线程池需要正值;多张表分别提交,不形成跨表事务。

iceberg.kafka.*

配置内部控制 Producer、Consumer 和 Admin 客户端。 此属性集合没有统一默认值;尝试读取 Worker properties 后用本前缀覆盖,具体属性的类型由对应 Kafka 客户端决定。
  • 类型Kafka 客户端透传属性
  • 默认值:N/A
  • 重要级别:N/A
  • 有效值 / 注意事项:自动读取依赖标准 Connect 启动入口和可读、包含 bootstrap.servers 的 Worker properties;不能依赖所有部署都能自动读取。必要时显式设置 Broker 地址和安全属性。Worker 的 consumer.*producer.* 不会自动去前缀。内部 Producer 强制生成的事务 ID 和序列化器;内部 Consumer 强制组 ID、关闭自动提交、使用 read_committed 和固定反序列化器,无显式设置时 auto.offset.reset 回退为 latest。此项不会覆盖 Task 数据消费者。

高级身份配置

iceberg.connect.group-id

告知协调器实际 Connect 数据消费组。
  • 类型string
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:未设置时读取回退为 connect-<connector name>。本项不会修改 Task 消费组;若通过 consumer.override.group.id 改变实际数据组,二者必须一致,否则可能无法正确选举协调器和完成提交。

iceberg.coordinator.transactional.prefix

设置生成事务 ID 时使用的前缀。
  • 类型string
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:读取时回退为空字符串;它不是完整事务 ID,不能用 iceberg.kafka.transactional.id 替换 Connector 生成的事务 ID。

最佳实践

将已知事件类别写入不同表

适用业务场景:同一事件流包含订单和支付事件,分析人员需要分别查询两类历史数据,而不是让每张表都接收全部事件。提前准备 analytics.ordersanalytics.payments,两张表分别具有对应事件的字段。 配置示例:在快速开始基础上覆盖目标表,并增加以下路由配置;保留其他 Catalog、Converter 和 Kafka 配置。输入值的 event_type 使用 orderpayment
关键说明:两个表达式进行全字符串匹配,每条有效类别记录只写入对应表。这不是按 Topic 自动映射表名。类别缺失、为空或不匹配时不会入表,消费位点仍可能推进;输入端须明确类别约定和无匹配数据的处理方式。若表达式重叠,同一记录会写入多张表,且两表快照不保证同时可见。

为新增事件流自动建表并接纳新增字段

适用业务场景:新的审计事件流尚无目标表,记录包含稳定的业务日期字段,之后会逐步增加查询所需字段。允许 Connector 从首条事件创建表,并在后续事件出现新增字段时扩展表结构,避免在每次新增字段前手工改表。 配置示例:在快速开始基础上,将目标表替换为尚不存在的 analytics.audit_events,增加自动建表、演进及日期字段分区配置。analytics 命名空间须可用,或 Catalog 允许创建。
输入端向快速开始指定的 Topic 写入不带 Schema 信封的 JSON 对象,首条记录包含用于分区的稳定字符串日期,后一条增加 actor 字段,例如:
关键说明:此处按 event_date 字段原值分区,不会把字符串自动识别成时间戳,也不会从写入时间生成业务日期。日期格式须由输入端保持一致。首条记录决定初始 Schema,后续非空字符串 actor 可作为新字段加入;这不代表能处理任意类型变化。分区规则只作用于新建表,须检查实际建表结果,因为无效规则可能退回无分区建表。已有表不会因开启自动建表而重建或重新分区。

监控

监控内容

关注 Kafka Connect 集群健康、Connector 与 Task 状态、吞吐、处理延迟、消费 Lag、Offset 提交、错误及重试,以及 Worker JVM 的堆内存、GC 和线程信号;仅在启用相应框架错误处理时关注 DLQ 活动。区分输入位点推进与目标表查询可见性,避免将低 Lag 直接解读为所有写入均已完成。

导入 Grafana 大盘

下载 Kafka Connect 集群 Grafana 大盘,确认监控采集已启用,Prometheus 兼容数据源及采集指标、标签与大盘查询匹配,然后在 Grafana 中导入 JSON 并选择对应数据源。

限制条件

  • 内置写入路径追加记录,不执行按主键的 upsert、唯一性约束或 CDC 行级删除;空值 Tombstone 不会删除已有行。
  • 静态目标表列表与动态路由模式互斥。
  • 多表快照独立提交,不提供跨表原子性或同时可见性。
  • Schema 演进不自动执行字段重命名、删列、类型收窄或任意破坏性变更。
  • 恢复依赖保留控制事件、相关消费组位点、引用的数据文件和快照检查点历史;清理这些资源或重置身份、位点可能破坏 exactly-once 恢复前提。
  • Connect Converter 和 SMT 阶段的容错及 DLQ 配置不能泛化为 Iceberg 写入路径中每条坏记录都能被跳过或送入 DLQ。

常见问题

消费 Lag 已经降低,为什么表里还看不到数据?

输入位点与文件引用先通过 Kafka 事务提交,表快照随后由协调器提交。先检查是否到达提交周期、目标表及查询分支是否正确,再查看协调器、Catalog 或存储错误。默认提交周期为 5 分钟,慢响应或 Catalog 提交会进一步影响可见性;不能以消费位点代替表提交结果。

Task 没有失败,为什么部分记录没有入表?

检查路由字段是否缺失、为空,或在静态路由中不满足全字符串正则。空值 Tombstone 也不会入表。动态路由关闭自动建表时,目标表不存在也会跳过写入;这不包括非法表标识符、权限错误或其他 Catalog 异常。这些跳过不会自动转入 DLQ,且后续 Kafka 事务成功时可提交对应消费位点。修正类别数据、匹配条件或补建动态目标表后,普通重启不会自动补回已提交位点对应的跳过记录,须另行确定补录范围和幂等处理策略。 静态路由关闭自动建表时,实际选中的目标缺表会使 Task 进入 FAILED 状态,不能按动态缺表跳过处理。检查异常中的目标表,核实命名空间、Schema、分区及权限,补建可加载且与输入兼容的目标资源后,保持 Connector 名、实际数据消费组、iceberg.connect.group-id、控制 Topic、控制消费组及其位点、目标表与提交分支和快照检查点身份不变,再重启失败 Task 继续处理。先核查失败输入对应的持久消费位点,不把内存中的位点推进当作已提交;重启后检查失败输入的重放及新输入的最终行集、快照检查点和持久位点,确认没有遗漏或重复。保留恢复所需的控制事件、引用的数据文件及检查点历史,不盲目换组或重置位点,也不将该处理方式视为任意故障下的多表原子恢复保证。

开启自动建表后,为什么表没有按预期分区?

分区配置只在创建缺失表时使用;已有表不会被重新分区。对新表,检查分区字段是否可推导、变换是否与字段类型兼容,以及日志是否有分区规则错误。发生解析或构建异常时可能创建无分区表,应按实际表结构处理,不能把继续写入理解为分区配置成功。

设置了标识列,为什么同一个业务键仍有多行?

标识列不启用唯一性约束或合并逻辑,内置路径会追加每条非空记录。CDC 信封规范化仅整理变更记录及操作元数据,不会把更新、删除自动应用到已有行。需要当前状态表时,须由下游查询逻辑或独立处理作业实现合并;exactly-once 也不会去重输入中本来就重复的业务事件。

修改了消费组后,为什么提交停滞?

检查 Task 的实际数据消费组是否与 iceberg.connect.group-id 一致。未覆盖时后者回退为 connect-<connector name>;仅修改 iceberg.connect.group-id 不会改变真实消费者。若配置了 consumer.override.group.id,须让两者一致并确认 Worker 允许该覆盖。消费组、控制 Topic 及控制组位点属于恢复身份,不应把任意更换或重置当作透明重启。

配置 DLQ 后,为什么坏记录仍使 Task 失败?

先区分异常来自 Converter、SMT,还是 Iceberg 记录转换、文件写入及 Catalog 操作。框架容错和 DLQ 仅覆盖其管理的阶段,不能承诺接管此 Connector 的 SinkTask.put 写入异常。根据完整异常和安全的 Topic、分区、位点上下文修正输入或目标 Schema;存储、事务及 Catalog 错误还须检查对应服务和权限,避免仅增加容错配置而丢失真实原因。

提交到其他分支后,为什么 main 查询没有变化?

检查全局与单表 commit-branch 的有效值,单表配置优先;未设置时才提交到 main。查询须指向实际写入分支。新分支不必提前创建,但已有同名引用若是 Tag 就不能作为提交目标;该配置不会把分支快照自动合并回 main