自定义主题自动创建
Kafka Connect 自动主题创建的自定义
Kafka 提供了两种自动创建主题的机制。你可以为 Kafka 代理(broker)启用自动主题创建;此外,从 Kafka 2.6.0 开始,你还可以让 Kafka Connect 创建主题。Kafka 代理使用 auto.create.topics.enable 属性来控制自动创建主题。在 Kafka Connect 中,topic.creation.enable 属性指明是否允许 Kafka Connect 创建主题。在这两种情况下,这些属性的默认设置都启用了自动创建主题。
启用自动创建主题后,如果 Debezium 源连接器为某个表发出变更事件记录,而该表尚不存在对应的目标主题,那么在事件记录被写入 Kafka 时,主题就会在运行时被创建。
代理端自动创建主题与 Kafka Connect 端自动创建主题的区别
由代理创建的主题只能共享一份默认配置。代理无法为不同的主题或主题集合应用各自不同的配置。相比之下,Kafka Connect 在创建主题时可以应用多种配置,按照 Debezium 连接器配置中的指定来设置副本因子、分区数量以及其他主题特定的设置。连接器配置定义了一组主题创建组(topic creation group),并将一组主题配置属性与每个组关联起来。
代理配置与 Kafka Connect 配置彼此独立。无论你是否在代理端禁用了主题创建,Kafka Connect 都可以创建主题。如果你同时在代理端和 Kafka Connect 中启用了自动创建主题,则以 Connect 的配置为准,只有当 Kafka Connect 配置中的设置都不适用时,代理才会创建主题。
为 Kafka 代理禁用自动主题创建
默认情况下,Kafka 代理的配置允许代理在主题不存在时于运行时创建主题。但是,代理无法为不同的连接器创建具有不同设置的主题。代理创建的所有主题都具有相同的保留期、日志压缩设置、分区数量、副本因子等。如果你有多个连接器,并且每个连接器都需要各自的主题配置,请使用 Kafka Connect 自动主题创建。
如果 Kafka Connect 的主题创建功能不可用(例如部署环境使用的是 Kafka 2.6 或更早版本),那么要保证自定义主题属性,唯一的办法就是自行预先创建主题,可以手动创建,也可以通过自定义的部署流程创建。
如前所述,代理端的主题创建与 Kafka Connect 的主题创建是相互独立的,因此要使用 Kafka Connect 机制,并不严格要求在代理端禁用主题创建。话虽如此,阻止代理创建主题有以下好处:
- 防止意外创建主题。
- 确保所有连接器主题都通过 Connect 配置创建。
- 使主题配置更可预测、更易于审计。
操作步骤
- 在代理配置中,将
auto.create.topics.enable的值设置为false。
设置 Kafka Connect
Kafka Connect 中的主题自动创建由 topic.creation.enable 属性控制。该属性的默认值为 true,即启用主题自动创建,如下例所示:
topic.creation.enable = truetopic.creation.enable 属性的设置适用于 Connect 集群中的所有工作进程。
Kafka Connect 的自动主题创建要求你定义 Kafka Connect 在创建主题时所应用的配置属性。你在 Debezium 连接器配置中通过定义主题组来指定主题配置属性,然后指定要应用于每个组的属性。连接器配置会定义一个默认主题创建组,另外(可选地)还可以定义一个或多个自定义主题创建组。自定义主题创建组使用主题名称模式列表来指定该组的设置所适用的主题。
有关 Kafka Connect 如何将主题与主题创建组匹配的详细信息,请参阅主题创建组。有关如何将配置属性分配给各组的更多信息,请参阅主题创建组配置属性。
默认情况下,Kafka Connect 创建的主题基于模式 server.schema.table 命名,例如 dbserver.myschema.inventory。
如果你不希望允许 Kafka Connect 自动创建主题,请在 Kafka Connect 配置(connect-distributed.properties 文件,或在使用 Debezium 的 Kafka Connect 容器镜像时通过环境变量 CONNECT_TOPIC_CREATION_ENABLE)中将 topic.creation.enable 的值设置为 false。 |
|---|
Kafka Connect 的自动主题创建要求至少为 default 主题创建组设置 replication.factor 和 partitions 属性。各组从 Kafka 代理的默认值中获取所需属性的值是有效的。 |
|---|
配置
要让 Kafka Connect 自动创建主题,它需要从源连接器获取在创建主题时要应用的配置属性信息。你在每个 Debezium 连接器的配置中定义控制主题创建的属性。当 Kafka Connect 为连接器发出的事件记录创建主题时,所生成的主题会从相应的分组获取其配置。该配置仅适用于该连接器发出的事件记录。
主题创建分组
一组主题属性与一个主题创建分组相关联。至少,你必须定义一个 default 主题创建分组并指定其配置属性。除此之外,你还可以选择定义一个或多个自定义主题创建分组,并为每个分组指定各自的属性。
在创建自定义主题创建分组时,你需要根据主题名称模式为每个分组定义成员主题。你可以指定命名模式,用来描述要包含在各分组中或从各分组中排除的主题。include 和 exclude 属性包含以逗号分隔的正则表达式列表,用于定义主题名称模式。例如,如果要让某个分组包含所有以字符串 dbserver1.inventory 开头的主题,请将其 topic.creation.inventory.include 属性的值设置为 dbserver1\\.inventory\\.*。
如果为自定义主题分组同时指定了 include 和 exclude 属性,则排除规则优先,并覆盖包含规则。 |
|---|
主题创建分组配置属性
default 主题创建分组以及每个自定义分组都关联着一组唯一的配置属性。你可以将分组配置为包含任意 Kafka 主题级配置属性。例如,你可以为某个主题分组指定旧主题段的清理策略、保留时间或主题压缩类型。你必须至少定义一组最小属性,以描述所创建主题的配置。
如果没有注册任何自定义分组,或者已注册分组的 include 模式与任何待创建主题的名称都不匹配,则 Kafka Connect 会使用 default 分组的配置来创建主题。
有关通用的主题配置注意事项,请参阅 Debezium 安装指南中的配置 Debezium 主题。
默认分组配置
在使用 Kafka Connect 的自动创建主题功能之前,必须先创建一个默认主题创建组,并为其定义配置。默认主题创建组的配置适用于名称不匹配任何自定义主题创建组 include 列表模式的主题。
操作步骤
要为
topic.creation.default组定义属性,请将其添加到连接器配置 JSON 中,如以下示例所示:{ ... "topic.creation.default.replication.factor": 3, (1) "topic.creation.default.partitions": 10, (2) "topic.creation.default.cleanup.policy": "compact", (3) "topic.creation.default.compression.type": "lz4" (4) ... }
你可以在 default 组的配置中包含任意 Kafka 主题级配置属性。
| 项目 | 说明 |
|---|---|
| 1 | topic.creation.default.replication.factor 定义默认组创建的主题的副本因子。replication.factor 对 default 组是必填项,对自定义组则是可选项。如果自定义组未设置该属性,将回退使用 default 组的值。使用 -1 表示采用 Kafka 代理的默认值。 |
| 2 | topic.creation.default.partitions 定义默认组创建的主题的分区数。partitions 对 default 组是必填项,对自定义组则是可选项。如果自定义组未设置该属性,将回退使用 default 组的值。使用 -1 表示采用 Kafka 代理的默认值。 |
| 3 | topic.creation.default.cleanup.policy 映射到主题级配置参数中的 cleanup.policy 属性,用于定义日志保留策略。 |
| 4 | topic.creation.default.compression.type 映射到主题级配置参数中的 compression.type 属性,用于定义消息在硬盘上的压缩方式。 |
表 1. default 主题创建组的连接器配置
自定义组仅在必填的 replication.factor 和 partitions 属性上回退使用 default 组的设置。如果自定义主题组的配置未定义其他属性,则不会应用 default 组中指定的值。 |
|---|
自定义组配置
你可以定义多个自定义主题组,每个组都有各自的配置。
操作步骤
若要定义自定义主题组,请在连接器 JSON 中添加
topic.creation.<group_name>.include属性,并在组名之后列出该自定义组的各项属性。以下示例展示了
inventory与applicationlogs两个自定义主题创建组的配置样例:{ ... (1) "topic.creation.inventory.include": "dbserver1\\.inventory\\.*", (2) "topic.creation.inventory.partitions": 20, "topic.creation.inventory.cleanup.policy": "compact", "topic.creation.inventory.delete.retention.ms": 7776000000, (3) "topic.creation.applicationlogs.include": "dbserver1\\.logs\\.applog-.*", (4) "topic.creation.applicationlogs.exclude": "dbserver1\\.logs\\.applog-old-.*", (5) "topic.creation.applicationlogs.replication.factor": 1, "topic.creation.applicationlogs.partitions": 20, "topic.creation.applicationlogs.cleanup.policy": "delete", "topic.creation.applicationlogs.retention.ms": 7776000000, "topic.creation.applicationlogs.compression.type": "lz4", ... }
| 项 | 说明 |
|---|---|
| 1 | 定义 inventory 组的配置。对于自定义组,replication.factor 和 partitions 属性是可选的。如果未设置值,自定义组将回退到为 default 组设置的值。将值设置为 -1 可使用为 Kafka 代理设置的值。 |
| 2 | topic.creation.inventory.include 定义了一个正则表达式,用于匹配所有以 dbserver1.inventory. 开头的主题。为 inventory 组定义的配置仅应用于名称与指定正则表达式匹配的主题。 |
| 3 | 定义 applicationlogs 组的配置。对于自定义组,replication.factor 和 partitions 属性是可选的。如果未设置值,自定义组将回退到为 default 组设置的值。将值设置为 -1 可使用为 Kafka 代理设置的值。 |
| 4 | topic.creation.applicationlogs.include 定义了一个正则表达式,用于匹配所有以 dbserver1.logs.applog- 开头的主题。为 applicationlogs 组定义的配置仅应用于名称与指定正则表达式匹配的主题。由于该组还定义了 exclude 属性,与 include 正则表达式匹配的主题可能还会受到该 exclude 属性的进一步限制。 |
| 5 | topic.creation.applicationlogs.exclude 定义了一个正则表达式,用于匹配所有以 dbserver1.logs.applog-old- 开头的主题。为 applicationlogs 组定义的配置仅应用于名称与给定正则表达式不匹配的主题。由于该组还定义了 include 属性,applicationlogs 组的配置仅应用于名称同时匹配指定的 include 正则表达式且不匹配指定的 exclude 正则表达式的主题。 |
表 2. 用于自定义 inventory 和 applicationlogs 主题创建组的连接器配置
注册自定义组
指定任何自定义主题创建组的配置后,注册这些组。
步骤
通过向连接器 JSON 中添加
topic.creation.groups属性,并指定以逗号分隔的组列表,来注册自定义组。以下示例注册了自定义主题创建组
inventory和applicationlogs:{ ... "topic.creation.groups": "inventory,applicationlogs", ... }
已完成的配置
下面的示例展示了一个完整的配置,其中既包含 default 主题组的配置,也包含 inventory 和 applicationlogs 这两个自定义主题创建组的配置:
示例:一个默认主题创建组与两个自定义组的配置
{
...
"topic.creation.default.replication.factor": 3,
"topic.creation.default.partitions": 10,
"topic.creation.default.cleanup.policy": "compact",
"topic.creation.default.compression.type": "lz4",
"topic.creation.groups": "inventory,applicationlogs",
"topic.creation.inventory.include": "dbserver1\\.inventory\\.*",
"topic.creation.inventory.partitions": 20,
"topic.creation.inventory.cleanup.policy": "compact",
"topic.creation.inventory.delete.retention.ms": 7776000000,
"topic.creation.applicationlogs.include": "dbserver1\\.logs\\.applog-.*",
"topic.creation.applicationlogs.exclude": "dbserver1\\.logs\\.applog-old-.*",
"topic.creation.applicationlogs.replication.factor": 1,
"topic.creation.applicationlogs.partitions": 20,
"topic.creation.applicationlogs.cleanup.policy": "delete",
"topic.creation.applicationlogs.retention.ms": 7776000000,
"topic.creation.applicationlogs.compression.type": "lz4"
}其他资源
有关主题自动创建的更多信息,可以参考以下资源:
- Debezium 博客:Auto-creating Debezium Change Data Topics
- 关于为 Kafka Connect 添加主题自动创建功能的 Kafka 改进提案:KIP-158 Kafka Connect should allow source connectors to set topic-specific settings for new topics
评论
登录后参与评论
KnowForge