分区路由
分区路由
默认情况下,当 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 表:
| id | name | description | weight |
|---|---|---|---|
| 101 | scooter | Small 2-wheel scooter | 3.14 |
| 102 | car battery | 12V car battery | 8.1 |
| 103 | 12-pack drill bits | 12-pack of drill bits with sizes ranging from devlive-community/knowforge#40 to devlive-community/knowforge#3 | 0.8 |
| 104 | hammer | 12oz carpenter’s hammer | 0.75 |
| 105 | hammer | 14oz carpenter’s hammer | 0.875 |
| 106 | hammer | 16oz carpenter’s hammer | 1.0 |
| 107 | rocks | box of assorted rocks | 5.3 |
| 108 | jacket | water resistent black wind breaker | 0.1 |
| 109 | spare tire | 24 inch spare tire | 22.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)配置选项
属性
默认值
说明
指定事件负载中 SMT 用于计算目标分区的字段。如果希望 SMT 将原始负载中的字段添加到输出数据结构的特定层级,可使用点号表示法。要访问与数据集合相关的字段,可以使用 after、before 或 change。change 是一个特殊字段,SMT 会根据操作类型自动在 after 或 before 元素中填充内容。如果记录中不存在指定的字段,SMT 会跳过它。例如:after.name,source.table,change.name
该 SMT 所作用的主题的分区数量。可使用 TopicNameMatches 谓词按主题过滤记录。
java
计算用于确定目标分区编号的字段哈希值时所使用的哈希函数。可设置以下值之一:
java
标准 Java Object::hashCode 函数
murmur
最新版本的 MurmurHash 函数,即 MurmurHash3
此属性为可选项。如果未指定值,或指定了无效值,则使用默认值。
评论
登录后参与评论
KnowForge