Kafka Connect
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 buildzip 压缩包位于 ./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 对象(
{...})。
- 它期望这些字符串中是 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 等的字段,其中名称表示元素在数组中的索引。每个元素随后作为对应字段的值传入。
评论
登录后参与评论
KnowForge