集成

Kafka Connect

师成师成· 更新于 2026-09-28· 阅读 33 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Kafka Connect

Kafka Connect 是一个流行的框架,通过连接器将数据移入和移出 Apache Kafka。它提供了许多不同的连接器,例如用于将数据从 Kafka 写入 S3 的 S3 sink 连接器,以及用于将关系型数据库中的变更数据捕获(CDC)记录写入 Kafka 的 Debezium source 连接器。

它采用简单、去中心化的分布式架构。集群由若干工作进程组成,连接器在这些进程上运行任务来执行工作。连接器的部署由配置驱动,因此通常无需编写任何代码即可运行连接器。

Apache Iceberg Sink Connector

Kafka Connect 的 Apache Iceberg Sink Connector 是一个 sink 连接器,用于将数据从 Kafka 写入 Iceberg 表。

功能特性

  • 针对集中式 Iceberg 提交的提交协调
  • 精确一次(exactly-once)交付语义
  • 多表分发(fan-out)
  • 自动建表与 schema 演进
  • 通过 Iceberg 的列映射功能实现字段名映射

安装

连接器的 zip 压缩包是 Iceberg 构建过程的一部分。你可以通过以下方式运行构建:

./gradlew -x test -x integrationTest clean build

zip 压缩包位于 ./kafka-connect/kafka-connect-runtime/build/distributions 目录下。其中一个发行版捆绑了 Hive Metastore 客户端及相关依赖,另一个则没有。请将发行版压缩包复制到所有节点的 Kafka Connect 插件目录中。

要求

该 Sink 依赖 KIP-447 来实现精确一次语义,这要求 Kafka 2.5 或更高版本。

配置

属性说明
iceberg.tables目标表的逗号分隔列表
iceberg.tables.dynamic-enabled设为 true 时,路由到 routeField 中指定的表,而不是使用 routeRegex,默认为 false
iceberg.tables.route-field用于多表扇出时,将记录路由到各表所使用的字段名称
iceberg.tables.default-commit-branch提交的默认分支,未指定时使用 main
iceberg.tables.default-id-columns默认的逗号分隔列列表,用于标识表中的行(主键)
iceberg.tables.default-partition-by创建表时使用的默认分区字段名的逗号分隔列表
iceberg.tables.auto-create-enabled设为 true 时自动创建目标表,默认为 false
iceberg.tables.evolve-schema-enabled设为 true 时将记录中缺失的字段添加到表结构中,默认为 false
iceberg.tables.schema-force-optional设为 true 时,在建表和演进表结构期间将列设置为可选;默认为 false,即遵循原有 schema
iceberg.tables.schema-case-insensitive设为 true 时按不区分大小写的名称查找表列,默认为 false,即区分大小写
iceberg.tables.auto-create-props.*自动建表时为新表设置的属性
iceberg.tables.write-props.*传递给 Iceberg 写入器初始化的属性,这些属性优先级更高
iceberg.table.<table-name>.commit-branch针对特定表的提交分支,未指定时使用 iceberg.tables.default-commit-branch
iceberg.table.<table-name>.id-columns逗号分隔的列列表,用于标识表中的行(主键)
iceberg.table.<table-name>.partition-by建表时使用的分区字段的逗号分隔列表
iceberg.table.<table-name>.route-regex用于将记录的 routeField 匹配到表的正则表达式
iceberg.control.topic控制主题的名称,默认为 control-iceberg
iceberg.control.group-id-prefix控制消费者组的前缀,默认为 cg-control
iceberg.control.commit.interval-ms提交间隔(毫秒),默认为 300,000(5 分钟)
iceberg.control.commit.timeout-ms提交超时时间(毫秒),默认为 30,000(30 秒)
iceberg.control.commit.threads提交所使用的线程数,默认为(cores * 2)
iceberg.coordinator.transactional.prefix协调器生产者所用事务 ID 的前缀,默认不使用前缀(空前缀)
iceberg.catalog目录(catalog)的名称,默认为 iceberg
iceberg.catalog.*传递给 Iceberg 目录初始化的属性
iceberg.hadoop-conf-dir若指定,则会加载该目录下的 Hadoop 配置文件
iceberg.hadoop.*传递给 Hadoop 配置的属性
iceberg.kafka.*传递给控制主题 Kafka 客户端初始化的属性

如果 iceberg.tables.dynamic-enabled 为 false(默认值),则必须指定 iceberg.tables。如果 iceberg.tables.dynamic-enabled 为 true,则必须指定 iceberg.tables.route-field,该字段将包含表的名称。

Kafka 配置

默认情况下,连接器会尝试使用 worker 属性中的 Kafka 客户端配置来连接控制主题。如果由于某种原因无法读取该配置,可以使用 iceberg.kafka.* 属性显式设置 Kafka 客户端设置。

消息格式

消息应使用相应的 Kafka Connect 转换器转换为结构体(struct)或映射(map)。

Catalog 配置

iceberg.catalog.* 属性是连接 Iceberg catalog 所必需的。默认发行版中包含了核心的 catalog 类型,包括 REST、Glue、DynamoDB、Hadoop、Nessie、JDBC、Hive 和 BigQuery Metastore。默认发行版不包含 JDBC 驱动,因此如有需要请自行添加。使用 Hive catalog 时,请使用包含 Hive metastore 客户端的发行版,否则需要自行包含该依赖。

要设置 catalog 类型,可以将 iceberg.catalog.type 设置为 rest、hive 或 hadoop。对于其他 catalog 类型,则需要改为将 iceberg.catalog.catalog-impl 设置为 catalog 类的名称。

REST 示例

"iceberg.catalog.type": "rest",
"iceberg.catalog.uri": "https://catalog-service",
"iceberg.catalog.credential": "<credential>",
"iceberg.catalog.warehouse": "<warehouse>",

Hive 示例

注意:请使用包含 HMS 客户端的发行版(或自行引入 HMS 客户端)。使用 S3 作为存储时请使用 S3FileIO,使用 GCS 时请使用 GCSFileIO(HiveCatalog 默认使用 HadoopFileIO)。

"iceberg.catalog.type": "hive",
"iceberg.catalog.uri": "thrift://hive:9083",
"iceberg.catalog.io-impl": "org.apache.iceberg.aws.s3.S3FileIO",
"iceberg.catalog.warehouse": "s3a://bucket/warehouse",
"iceberg.catalog.client.region": "us-east-1",
"iceberg.catalog.s3.access-key-id": "<AWS access>",
"iceberg.catalog.s3.secret-access-key": "<AWS secret>",

Glue 示例

"iceberg.catalog.catalog-impl": "org.apache.iceberg.aws.glue.GlueCatalog",
"iceberg.catalog.warehouse": "s3a://bucket/warehouse",
"iceberg.catalog.io-impl": "org.apache.iceberg.aws.s3.S3FileIO",

Nessie 示例

"iceberg.catalog.catalog-impl": "org.apache.iceberg.nessie.NessieCatalog",
"iceberg.catalog.uri": "http://localhost:19120/api/v2",
"iceberg.catalog.ref": "main",
"iceberg.catalog.warehouse": "s3a://bucket/warehouse",
"iceberg.catalog.io-impl": "org.apache.iceberg.aws.s3.S3FileIO",

BigQuery 元数据存储示例

"iceberg.catalog.catalog-impl": "org.apache.iceberg.gcp.bigquery.BigQueryMetastoreCatalog",
"iceberg.catalog.gcp.bigquery.project-id": "my-project",
"iceberg.catalog.gcp.bigquery.location": "us-east1",
"iceberg.catalog.warehouse": "gs://bucket/warehouse",
"iceberg.catalog.io-impl": "org.apache.iceberg.gcp.gcs.GCSFileIO",
"iceberg.tables.auto-create-props.bq_connection": "projects/my-project/locations/us-east1/connections/my-connection",

注意

根据你的配置,可能还需要设置 iceberg.catalog.s3.endpoint、iceberg.catalog.s3.staging-dir 或 iceberg.catalog.s3.path-style-access。有关配置目录的完整详情,请参阅 Iceberg 文档。

Azure ADLS 配置示例

使用 ADLS 时,Azure 要求为其 Java SDK 传入 AZURE_CLIENT_ID、AZURE_TENANT_ID 和 AZURE_CLIENT_SECRET。如果你在容器中运行 Kafka Connect,请务必将这些值作为环境变量注入。更多信息请参阅 Azure Identity Java 客户端库。

示例如下:

AZURE_CLIENT_ID=e564f687-7b89-4b48-80b8-111111111111
AZURE_TENANT_ID=95f2f365-f5b7-44b1-88a1-111111111111
AZURE_CLIENT_SECRET="XXX"

其中 CLIENT_ID 是应用注册中已注册应用的应用 ID,TENANT_ID 来自你的 Azure 租户属性,CLIENT_SECRET 则在选择你的具体应用注册后,于"管理"下的"证书和密码"部分中创建。你可能需要在中间面板中选择"客户端密码",然后点击"新客户端密码"前的"+"来生成一个。请务必确保将该变量设置为 Value(值),而不是 Id(标识符)。

同样重要的是,该应用注册必须在你的存储账户的访问控制(IAM)中被授予"Storage Blob Data Contributor"角色分配,否则它将无法在那里写入新文件。

接下来,在 Connector 的配置中,你需要包含以下内容:

"iceberg.catalog.type": "rest",
"iceberg.catalog.uri": "https://catalog:8181",
"iceberg.catalog.warehouse": "abfss://storage-container-name@storageaccount.dfs.core.windows.net/warehouse",
"iceberg.catalog.io-impl": "org.apache.iceberg.azure.adlsv2.ADLSFileIO",
"iceberg.catalog.include-credentials": "true"

其中 storage-container-name 是 Azure 存储账户中的容器名称,/warehouse 是该容器中默认写入 Apache Iceberg 文件的位置(或当 iceberg.tables.auto-create-enabled=true 时使用),include-credentials 参数会一并传递 Azure Java 客户端凭据。这样配置后,Iceberg Sink 连接器将连接到 iceberg.catalog.uri 指定的 REST catalog 实现,以获取 ADLSv2 客户端所需的连接字符串。

Google GCS 配置示例

默认情况下,将使用应用默认凭据(ADC)连接到 GCS。有关 ADC 工作原理的详细信息,请参阅 Google Cloud 文档。

"iceberg.catalog.type": "rest",
"iceberg.catalog.uri": "https://catalog:8181",
"iceberg.catalog.warehouse": "gs://bucket-name/warehouse",
"iceberg.catalog.io-impl": "org.apache.iceberg.gcp.gcs.GCSFileIO"

Hadoop 配置

使用 HDFS 或 Hive 时,sink 会初始化 Hadoop 配置。首先,从 classpath 加载配置文件。接着,如果指定了 iceberg.hadoop-conf-dir,则从该位置加载配置文件。最后,应用 sink 配置中的所有 iceberg.hadoop.* 属性。合并这些配置时,优先级顺序为:sink 配置 > 配置目录 > classpath。

示例

初始设置

源主题

这里假设源主题已存在,名称为 events。

控制主题

如果你的 Kafka 集群将 auto.create.topics.enable 设置为 true(默认值),则控制主题会被自动创建。否则,你需要先手动创建该主题。默认的主题名称为 control-iceberg:

bin/kafka-topics.sh  \
  --command-config command-config.props \
  --bootstrap-server ${CONNECT_BOOTSTRAP_SERVERS} \
  --create \
  --topic control-iceberg \
  --partitions 1

注意:运行在 Confluent Cloud 上的集群默认将 auto.create.topics.enable 设置为 false。

Iceberg catalog 配置

带有 iceberg.catalog. 前缀的配置属性会传递给 Iceberg catalog 初始化。有关如何配置特定 catalog 的详细信息,请参阅 Iceberg 文档。

单个目标表

此示例将所有传入记录写入单个表。

创建目标表

CREATE TABLE default.events (
    id STRING,
    type STRING,
    ts TIMESTAMP,
    payload STRING)
PARTITIONED BY (hours(ts))

连接器配置

此示例配置连接到 Iceberg REST 目录。

{
  "name": "events-sink",
  "config": {
    "connector.class": "org.apache.iceberg.connect.IcebergSinkConnector",
    "tasks.max": "2",
    "topics": "events",
    "iceberg.tables": "default.events",
    "iceberg.catalog.type": "rest",
    "iceberg.catalog.uri": "https://localhost",
    "iceberg.catalog.credential": "<credential>",
    "iceberg.catalog.warehouse": "<warehouse name>"
  }
}

多表扇出、静态路由

此示例会将 type 设置为 list 的记录写入表 default.events_list,并将 type 设置为 create 的记录写入表 default.events_create。其他记录将被跳过。

创建两张目标表

CREATE TABLE default.events_list (
    id STRING,
    type STRING,
    ts TIMESTAMP,
    payload STRING)
PARTITIONED BY (hours(ts));

CREATE TABLE default.events_create (
    id STRING,
    type STRING,
    ts TIMESTAMP,
    payload STRING)
PARTITIONED BY (hours(ts));

连接器配置

{
  "name": "events-sink",
  "config": {
    "connector.class": "org.apache.iceberg.connect.IcebergSinkConnector",
    "tasks.max": "2",
    "topics": "events",
    "iceberg.tables": "default.events_list,default.events_create",
    "iceberg.tables.route-field": "type",
    "iceberg.table.default.events_list.route-regex": "list",
    "iceberg.table.default.events_create.route-regex": "create",
    "iceberg.catalog.type": "rest",
    "iceberg.catalog.uri": "https://localhost",
    "iceberg.catalog.credential": "<credential>",
    "iceberg.catalog.warehouse": "<warehouse name>"
  }
}

多表扇出、动态路由

本示例根据 db_table 字段中的值写入对应的表。如果该名称的表不存在,则该条记录会被跳过。例如,如果记录的 db_table 字段设置为 default.events_list,则该记录会被写入 default.events_list 表。

创建两个目标表

创建两张表的方法见上文。

连接器配置

{
  "name": "events-sink",
  "config": {
    "connector.class": "org.apache.iceberg.connect.IcebergSinkConnector",
    "tasks.max": "2",
    "topics": "events",
    "iceberg.tables.dynamic-enabled": "true",
    "iceberg.tables.route-field": "db_table",
    "iceberg.catalog.type": "rest",
    "iceberg.catalog.uri": "https://localhost",
    "iceberg.catalog.credential": "<credential>",
    "iceberg.catalog.warehouse": "<warehouse name>"
  }
}

用于 Apache Iceberg Sink Connector 的 SMT

本项目包含一些 SMT(简单消息转换器),在转换 Kafka 数据供 Iceberg sink connector 使用时可能会很有用。

CopyValue

(实验性功能)

CopyValue SMT 将一个字段的值复制到一个新的字段中。

配置

属性说明
source.field源字段名称
target.field目标字段名称

示例

"transforms": "copyId",
"transforms.copyId.type": "org.apache.iceberg.connect.transforms.CopyValue",
"transforms.copyId.source.field": "id",
"transforms.copyId.target.field": "id_copy",

DmsTransform

(实验性)

DmsTransform SMT 会转换 AWS DMS 格式的消息,供接收端的 CDC 功能使用。它会将 data 元素中的字段提升到顶层,并添加以下元数据字段:_cdc.op、_cdc.ts 和 _cdc.source。

配置

该 SMT 目前没有可配置项。

DebeziumTransform

(实验性)

DebeziumTransform SMT 会转换 Debezium 格式的消息,供接收端的 CDC 功能使用。它会将 before 或 after 元素中的字段提升到顶层,并添加以下元数据字段:_cdc.op、_cdc.ts、_cdc.offset、_cdc.source、_cdc.target 和 _cdc.key。

配置
属性说明
cdc.target.pattern用于设置 CDC 目标字段值的模式,默认为 {db}.{table}

JsonToMapTransform

(实验性)

JsonToMapTransform SMT 会将字符串解析为 JSON 对象负载,以推断 Schema。iceberg-kafka-connect 连接器针对无 Schema 数据(例如 Kafka 自带的 JsonConverter 产生的 Map)的目标是将 Map 转换为 Iceberg Struct。当 JSON 结构良好时这没有问题,但当 JSON 对象的键动态变化时,会导致 Iceberg 表因 Schema 演进而产生大量列。

在 JSON 结构不够规整的场景下,这个 SMT 非常有用,它能先把数据写入 Iceberg,再由查询引擎将其处理成更易于管理的形式。它会将嵌套对象转换为 Map,并在 Schema 中包含 Map 类型。连接器会遵循该 Schema,为 JSON 对象创建带有 Iceberg Map(String)列的 Iceberg 表。

注意:

  • 连接器的 value.converter 设置必须使用 stringConverter,而不是 jsonConverter

    • 它期望这些字符串中是 JSON 对象({...})。
  • 消息键、墓碑消息和消息头不会被转换,SMT 会原样传递。

配置
属性说明(默认值)
json.root(false)布尔值,表示从根节点开始处理

transforms.IDENTIFIER_HERE.json.root 适用于最不规整的数据。它会构造一个只包含单个名为 payload 字段的 Struct,其 Schema 为 Map<String, String>。

如果 transforms.IDENTIFIER_HERE.json.root 为 false(默认值),则会构造一个 Struct,其基本类型和数组字段的 Schema 由推断得出。嵌套对象会成为类型为 Map<String, String> 的字段。

包含空数组或空对象的键会被从最终 Schema 中过滤掉。数组会根据 JSON 数组的实际类型确定类型;如果 JSON 数组中包含混合类型,则会被转换为字符串数组。

示例 JSON:

{
  "key": 1,
  "array": [1,"two",3],
  "empty_obj": {},
  "nested_obj": {"some_key": ["one", "two"]}
}

如果 json.root 为 true,则结果如下:

SinkRecord.schema:
  "payload" : (Optional) Map<String, String>

Sinkrecord.value (Struct):
  "payload"  : Map(
    "key" : "1",
    "array" : "[1,"two",3]"
    "empty_obj": "{}"
    "nested_obj": "{"some_key":["one","two"]}"
  )

如果 json.root 为 false,则将变成以下内容

SinkRecord.schema:
  "key": (Optional) Int32,
  "array": (Optional) Array<String>,
  "nested_object": (Optional) Map<string, String>

SinkRecord.value (Struct):
  "key" 1,
  "array" ["1", "two", "3"]
  "nested_object" Map ("some_key" : "["one", "two"]")

KafkaMetadataTransform

(实验性功能)

KafkaMetadata 会注入 topic、partition、offset、timestamp,这些是 Kafka 消息的属性。

配置

属性说明(默认值)
field_name(_kafka_metadata)字段前缀
nested(false)为 true 时将数据嵌套到一个 struct 中,否则以带前缀的字段形式添加到顶层
external_field(无)向元数据追加一个常量 key,value(例如集群名称)

当 nested 开启时:

_kafka_metadata.topic、_kafka_metadata.partition、_kafka_metadata.offset、_kafka_metadata.timestamp

当 nested 关闭时:_kafka_metadata_topic、_kafka_metadata_partition、_kafka_metadata_offset、_kafka_metadata_timestamp

MongoDebeziumTransform

(实验性功能)

MongoDebeziumTransform SMT 会将带有 before/after BSON 字符串的 Mongo Debezium 格式消息,转换为 DebeziumTransform SMT 所期望的、带有类型的 before/after Struct。

它(目前)不支持在 mongodb 列不被底层目录类型支持时重命名列。

配置

属性说明
array_handling_mode设为 array 或 document,用于设置数组处理模式

值为 array(默认)时,会将数组编码为 array 数据类型。用户需要自行确保给定数组实例中的所有元素类型一致。这个选项有较多限制,但便于下游客户端处理数组。

值为 document 时,会以类似 BSON 序列化的方式将数组转换为 struct 的 struct。主 struct 包含名为 _0、_1、_2 等的字段,其中名称表示元素在数组中的索引。每个元素随后作为对应字段的值传入。

评论

登录后参与评论

正在加载评论…