变形

分区路由

qianmoQqianmoQ· 更新于 2026-09-28· 阅读 17 分钟· 0 次阅读

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

分区路由

默认情况下,当 Debezium 检测到某个数据集合中的变更时,它发出的变更事件会被发送到使用单个 Apache Kafka 分区的主题中。如自定义 Kafka Connect 自动创建主题所述,你可以自定义默认配置,基于主键的哈希值将事件路由到多个分区。

不过,在某些情况下,你可能还希望 Debezium 将事件路由到特定的主题分区。分区路由 SMT 允许你根据一个或多个指定的有效载荷(payload)字段的值,将事件路由到特定的目标分区。为计算目标分区,Debezium 会使用指定字段值的哈希。

示例:基本配置

你在 Debezium 连接器的 Kafka Connect 配置中配置分区路由转换。该配置指定以下参数:

partition.payload.fields

指定事件有效载荷中 SMT 用于计算目标分区的字段。你可以使用点号表示法来指定嵌套的有效载荷字段。

partition.topic.num

指定目标主题中的分区数量。

partition.hash.function

指定用于对字段求哈希的哈希函数,以确定目标分区的编号。

默认情况下,Debezium 会将已配置数据集合的所有变更事件记录发送到单个 Apache Kafka 主题。连接器不会将事件记录定向到该主题中的特定分区。

要配置 Debezium 连接器将事件路由到特定分区,需在 Debezium 连接器的 Kafka Connect 配置中配置 PartitionRouting SMT。

例如,你可以在连接器配置中添加如下配置。

...
topic.creation.default.partitions=2
topic.creation.default.replication.factor=1
...

topic.prefix=fulfillment
transforms=PartitionRouting
transforms.PartitionRouting.type=io.debezium.transforms.partitions.PartitionRouting
transforms.PartitionRouting.partition.payload.fields=change.name
transforms.PartitionRouting.partition.topic.num=2
transforms.PartitionRouting.predicate=allTopic
predicates=allTopic
predicates.allTopic.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.allTopic.pattern=fulfillment.*
...

根据上述配置,每当 SMT 接收到一条目标主题名称以 fulfillment 前缀开头的消息时,它会将该消息重定向到特定的主题分区。

SMT 通过消息负载中 name 字段的值计算目标分区。通过指定 allTopic 谓词,配置可以有选择性地应用 SMT。change 前缀是一个特殊关键字,它使 SMT 能够自动引用负载中描述数据 before(之前)或 after(之后)状态的元素。如果指定的字段不存在于事件消息中,SMT 将忽略它。如果消息中不存在任何字段,则该转换会完全忽略事件消息,并将原始版本的消息发送到默认目标主题。SMT 配置中由 topic.num 设置指定的分区数必须与 Kafka Connect 配置中指定的分区数一致。例如,在上面的配置示例中,Kafka Connect 属性 topic.creation.default.partitions 指定的值与 SMT 配置中的 topic.num 值相匹配。

给定如下 Products 表:

idnamedescriptionweight
101scooterSmall 2-wheel scooter3.14
102car battery12V car battery8.1
10312-pack drill bits12-pack of drill bits with sizes ranging from devlive-community/knowforge#40 to devlive-community/knowforge#30.8
104hammer12oz carpenter’s hammer0.75
105hammer14oz carpenter’s hammer0.875
106hammer16oz carpenter’s hammer1.0
107rocksbox of assorted rocks5.3
108jacketwater resistent black wind breaker0.1
109spare tire24 inch spare tire22.2

表 1. Products 表

根据该配置,SMT 将字段名称为 hammer 的记录的变更事件路由到同一个分区。也就是说,id 值为 104、105 和 106 的条目将被路由到同一个分区。

示例:高级配置

假设你想将来自两个数据集合(t1、t2)的事件路由到同一个主题(例如 my_topic),并且你希望使用字段 f1 对数据集合 t1 的事件进行分区,使用字段 f2 对数据集合 t2 的事件进行分区。

你可以应用如下配置:

transforms=PartitionRouting
transforms.PartitionRouting.type=io.debezium.transforms.partitions.PartitionRouting
transforms.PartitionRouting.partition.payload.fields=change.f1,change.f2
transforms.PartitionRouting.partition.topic.num=2
transforms.PartitionRouting.predicate=myTopic

predicates=myTopic
predicates.myTopic.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.myTopic.pattern=my_topic

上述配置并未指定如何重新路由事件,使其发送到特定的目标主题。有关如何将事件发送到默认目标主题以外的主题,请参阅 Topic Routing SMT。有关如何将事件发送到默认目标主题以外的主题,请参阅 Topic Routing SMT。

从 Debezium ComputePartition SMT 迁移

Debezium 的 ComputePartition SMT 已停止维护。以下章节介绍如何从 ComputePartition SMT 迁移到新的 PartitionRouting SMT。

假设配置为所有主题设置了相同的分区数,请将原来的 ComputePartition` 配置替换为 `PartitionRouting SMT。以下示例对比了这两种配置。

示例:旧版 ComputePartition 配置

...
topic.creation.default.partitions=2
topic.creation.default.replication.factor=1
...
topic.prefix=fulfillment
transforms=ComputePartition
transforms.ComputePartition.type=io.debezium.transforms.partitions.ComputePartition
transforms.ComputePartition.partition.data-collections.field.mappings=inventory.products:name,inventory.orders:purchaser
transforms.ComputePartition.partition.data-collections.partition.num.mappings=inventory.products:2,inventory.orders:2
...

将前文的 ComputePartition 替换为以下 PartitionRouting 配置。示例:替换先前 ComputePartition 配置的 PartitionRouting 配置

...
topic.creation.default.partitions=2
topic.creation.default.replication.factor=1
...

topic.prefix=fulfillment
transforms=PartitionRouting
transforms.PartitionRouting.type=io.debezium.transforms.partitions.PartitionRouting
transforms.PartitionRouting.partition.payload.fields=change.name,change.purchaser
transforms.PartitionRouting.partition.topic.num=2
transforms.PartitionRouting.predicate=allTopic
predicates=allTopic
predicates.allTopic.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.allTopic.pattern=fulfillment.*
...

如果 SMT 将事件输出到分区数不相同的多个主题,则必须为每个主题指定唯一的 partition.num.mappings 值。例如,在下面的示例中,旧版 products 集合对应的主题配置为 3 个分区,orders 数据集合对应的主题配置为 2 个分区:

示例:为不同主题设置唯一分区值的旧版 ComputePartition 配置

...
topic.prefix=fulfillment
transforms=ComputePartition
transforms.ComputePartition.type=io.debezium.transforms.partitions.ComputePartition
transforms.ComputePartition.partition.data-collections.field.mappings=inventory.products:name,inventory.orders:purchaser
transforms.ComputePartition.partition.data-collections.partition.num.mappings=inventory.products:3,inventory.orders:2
...

将前面的 ComputePartition 配置替换为下面的 PartitionRouting 配置: .PartitionRouting 配置,为不同的主题设置唯一的 partition.topic.num 值

...
topic.prefix=fulfillment

transforms=ProductsPartitionRouting,OrdersPartitionRouting
transforms.ProductsPartitionRouting.type=io.debezium.transforms.partitions.PartitionRouting
transforms.ProductsPartitionRouting.partition.payload.fields=change.name
transforms.ProductsPartitionRouting.partition.topic.num=3
transforms.ProductsPartitionRouting.predicate=products

transforms.OrdersPartitionRouting.type=io.debezium.transforms.partitions.PartitionRouting
transforms.OrdersPartitionRouting.partition.payload.fields=change.purchaser
transforms.OrdersPartitionRouting.partition.topic.num=2
transforms.OrdersPartitionRouting.predicate=products

predicates=products,orders
predicates.products.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.products.pattern=fulfillment.inventory.products
predicates.orders.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.orders.pattern=fulfillment.inventory.orders
...

配置选项

下表列出了可为分区路由 SMT 设置的配置选项。

表 2. 分区路由 SMT(PartitionRouting)配置选项

属性

默认值

说明

partition.payload.fields

指定事件负载中 SMT 用于计算目标分区的字段。如果希望 SMT 将原始负载中的字段添加到输出数据结构的特定层级,可使用点号表示法。要访问与数据集合相关的字段,可以使用 after、before 或 change。change 是一个特殊字段,SMT 会根据操作类型自动在 after 或 before 元素中填充内容。如果记录中不存在指定的字段,SMT 会跳过它。例如:after.name,source.table,change.name

partition.topic.num

该 SMT 所作用的主题的分区数量。可使用 TopicNameMatches 谓词按主题过滤记录。

partition.hash.function

java

计算用于确定目标分区编号的字段哈希值时所使用的哈希函数。可设置以下值之一:

java

标准 Java Object::hashCode 函数

murmur

最新版本的 MurmurHash 函数,即 MurmurHash3

此属性为可选项。如果未指定值,或指定了无效值,则使用默认值。

评论

登录后参与评论

正在加载评论…