概述
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 对象追加到已有表。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 处理后的非空顶层值须为
Struct或Map;转换器子属性须与实际输入编码匹配。
consumer.override.*
覆盖 Connect Task 的数据消费者属性。
此属性集合没有统一默认值,具体属性的类型和默认值由 Kafka Consumer 决定。
- 类型:
Kafka Consumer透传属性 - 默认值:N/A
- 重要级别:N/A
- 有效值 / 注意事项:须通过 Worker 的客户端覆盖策略;优先级为 Connect 基值、Worker
consumer.*、Connectorconsumer.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
- 必填:是
- 有效值 / 注意事项:
type与catalog-impl互斥。type不区分大小写,可选择hive、hadoop、rest、glue、nessie、jdbc;实际可用性取决于已安装的实现和客户端依赖。两者都未设置时底层回退到 Hive,但这不是 Connector 的默认类型声明。uri、warehouse、io-impl等是否必填由选定实现决定。
iceberg.hadoop-conf-dir
指定 Hadoop XML 配置目录。
- 类型:
string - 默认值:
null - 重要级别:中
- 有效值 / 注意事项:目录须能被 Worker 访问,读取其中的
core-site.xml、hdfs-site.xml和hive-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 - 重要级别:中
- 有效值 / 注意事项:动态模式必填;支持点分嵌套
Struct或Map字段路径。静态模式下,该字段的值与每张表的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还支持int到long、float到double的提升和必填字段放宽为可选。无 Schema 数据不具有同样的已有字段提升路径;关闭时额外字段会被忽略。它不支持任意重命名、删列或类型迁移。
iceberg.tables.schema-force-optional
将创建或演进的字段设为可选。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:也作用于创建或演进的容器值类型;不会整体改写已有表所有字段的可选性。默认尊重输入 Schema 的可选性。
iceberg.tables.schema-case-insensitive
控制表字段查找是否忽略大小写。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:不重命名字段;表存在名称映射时使用映射查找路径,此开关不会无条件覆盖映射。
iceberg.tables.default-partition-by
设置自动建表的默认分区规则。
- 类型:
string - 默认值:
null - 重要级别:中
- 有效值 / 注意事项:仅作用于新建表,不改变已有表分区。支持字段本身以及
year、month、day、hour、bucket、truncate变换;多个规则用逗号分隔,变换括号内的逗号保留,bucket和truncate的参数顺序为列名、宽度。解析或构建异常时会记录错误并退回无分区建表,不能只凭配置判断分区已生效。
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.orders 和 analytics.payments,两张表分别具有对应事件的字段。
配置示例:在快速开始基础上覆盖目标表,并增加以下路由配置;保留其他 Catalog、Converter 和 Kafka 配置。输入值的 event_type 使用 order 或 payment。
为新增事件流自动建表并接纳新增字段
适用业务场景:新的审计事件流尚无目标表,记录包含稳定的业务日期字段,之后会逐步增加查询所需字段。允许 Connector 从首条事件创建表,并在后续事件出现新增字段时扩展表结构,避免在每次新增字段前手工改表。 配置示例:在快速开始基础上,将目标表替换为尚不存在的analytics.audit_events,增加自动建表、演进及日期字段分区配置。analytics 命名空间须可用,或 Catalog 允许创建。
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。