变形

MongoDB Outbox 事件路由器

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

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

MongoDB Outbox Event Router

发件箱模式(outbox pattern)提供了一种在多个微服务之间安全、可靠地交换数据的方式。发件箱模式的实现可确保服务存储在其数据库中的数据,与 Debezium 发布给其他服务消费的数据保持一致。

此 SMT 仅适用于 Debezium MongoDB 连接器。有关在关系型数据库中使用 outbox event router SMT 的信息,请参阅 Outbox event router。

要在 Debezium 应用程序中实现发件箱模式,请将 Debezium 连接器配置为:

  • 捕获 outbox 集合中的更改
  • 应用 Debezium MongoDB outbox 事件路由器单消息转换(SMT)

配置为应用 MongoDB outbox SMT 的 Debezium 连接器,应当只捕获 outbox 集合中发生的更改。更多信息请参阅选择性应用转换的选项。

连接器可以捕获多个 outbox 集合中的更改,但前提是每个 outbox 集合具有相同的结构。

要使用此 SMT,对实际业务集合的操作以及向 outbox 集合的插入必须作为多文档事务的一部分来执行(MongoDB 自 4.0 起支持多文档事务),以防止业务集合与 outbox 集合之间出现潜在的数据不一致。为了未来的更新,在不使用多文档事务的情况下,能够在 ACID 事务中更新现有数据并插入 outbox 事件,我们计划支持额外的配置,将 outbox 事件以现有集合的子文档形式存储,而不是存放在独立的 outbox 集合中。

有关发件箱模式的更多信息,请参阅 Reliable Microservices Data Exchange With the Outbox Pattern。

示例 outbox 消息

为了理解如何配置 Debezium MongoDB outbox 事件路由器 SMT,请看下面这个 Debezium outbox 消息的示例:

# Kafka Topic: outbox.event.order
# Kafka Message key: "b2730779e1f596e275826f08"
# Kafka Message Headers: "id=596e275826f08b2730779e1f"
# Kafka Message Timestamp: 1556890294484
{
  "{\"id\": {\"$oid\": \"da8d6de63b7745ff8f4457db\"}, \"lineItems\": [{\"id\": 1, \"item\": \"Debezium in Action\", \"status\": \"ENTERED\", \"quantity\": 2, \"totalPrice\": 39.98}, {\"id\": 2, \"item\": \"Debezium for Dummies\", \"status\": \"ENTERED\", \"quantity\": 1, \"totalPrice\": 29.99}], \"orderDate\": \"2019-01-31T12:13:01\", \"customerId\": 123}"
}

配置了 MongoDB 发件箱事件路由器 SMT 的 Debezium 连接器,会通过转换原始 Debezium 变更事件消息来生成上面的消息,转换示例如下:

# Kafka Message key: { "id": "{\"$oid\": \"596e275826f08b2730779e1f\"}" }
# Kafka Message Headers: ""
# Kafka Message Timestamp: 1556890294484
{
  "patch": null,
  "after": "{\"_id\": {\"$oid\": \"596e275826f08b2730779e1f\"}, \"aggregateid\": {\"$oid\": \"b2730779e1f596e275826f08\"}, \"aggregatetype\": \"Order\", \"type\": \"OrderCreated\", \"payload\": {\"_id\": {\"$oid\": \"da8d6de63b7745ff8f4457db\"}, \"lineItems\": [{\"id\": 1, \"item\": \"Debezium in Action\", \"status\": \"ENTERED\", \"quantity\": 2, \"totalPrice\": 39.98}, {\"id\": 2, \"item\": \"Debezium for Dummies\", \"status\": \"ENTERED\", \"quantity\": 1, \"totalPrice\": 29.99}], \"orderDate\": \"2019-01-31T12:13:01\", \"customerId\": 123}}",
  "source": {
    "version": "3.6.3.Final",
    "connector": "mongodb",
    "name": "fulfillment",
    "ts_ms": 1558965508000,
    "ts_us": 1558965508000000,
    "ts_ns": 1558965508000000000,
    "snapshot": false,
    "db": "inventory",
    "rs": "rs0",
    "collection": "customers",
    "ord": 31,
    "h": 1546547425148721999
  },
  "op": "c",
  "ts_ms": 1556890294484,
  "ts_us": 1556890294484452,
  "ts_ns": 1556890294484452697,
}

此 Debezium 发件箱消息示例基于默认发件箱事件路由配置,该配置假定发件箱集合的结构以及基于聚合的事件路由方式。若要自定义行为,发件箱事件路由 SMT 提供了众多配置选项。

基本发件箱集合

要应用默认的 MongoDB 发件箱事件路由 SMT 配置,你的发件箱集合应包含以下字段:

{
  "_id": "objectId",
  "aggregatetype": "string",
  "aggregateid": "objectId",
  "type": "string",
  "payload": "object"
}

表 1. 预期发件箱集合字段说明

字段 作用

id

包含事件的唯一 ID。在发件箱消息中,该值是一个消息头。你可以使用该 ID,例如,来删除重复消息。

要从其他发件箱集合字段获取事件的唯一 ID,请在连接器配置中设置 collection.field.event.id SMT 选项。

aggregatetype

包含一个值,SMT 会将其追加到连接器发出发件箱消息所用的主题名称之后。默认情况下,该值会替换 route.topic.replacement SMT 选项中默认的 ${routedByValue} 变量。

例如,在默认配置中,route.by.field SMT 选项被设置为 aggregatetype,而 route.topic.replacement SMT 选项被设置为 outbox.event.${routedByValue}。

假设你的应用向发件箱集合中添加了两个文档。第一个文档中 aggregatetype 字段的值为 customers,第二个文档中 aggregatetype 字段的值为 orders。连接器会将第一个文档发出到 outbox.event.customers 主题,将第二个文档发出到 outbox.event.orders 主题。

要从其他发件箱集合字段获取该值,请在连接器配置中设置 route.by.field SMT 选项。

aggregateid

包含事件键,它为有效负载提供一个 ID。SMT 将该值用作所发出的发件箱消息的键。这对于在 Kafka 分区中保持正确的顺序非常重要。

要从其他发件箱集合字段获取事件键,请在连接器配置中设置 collection.field.event.key SMT 选项。

payload

发件箱变更事件的表示形式。默认结构为 JSON。默认情况下,Kafka 消息值仅由 payload 的值组成。但是,如果配置了发件箱事件以包含额外字段,则 Kafka 消息值包含一个封装有效负载和这些额外字段的信封,并且每个字段分别表示。有关更多信息,请参阅发出包含额外字段的消息。

要从发件箱集合的其他字段获取事件负载,请在连接器配置中设置 collection.field.event.payload SMT 选项。

附加自定义字段

发件箱集合中的任何其他字段都可以添加到发件箱事件中,既可以放在负载部分,也可以作为消息标头。

一个例子是 eventType 字段,它承载用户自定义的值,用于帮助分类或组织事件。

基本配置

要配置 Debezium MongoDB 连接器以支持发件箱模式,请配置 outbox.MongoEventRouter SMT。要获得该 SMT 的默认行为,只需将其添加到连接器配置中而不指定任何选项,如以下示例所示:

transforms=outbox,...
transforms.outbox.type=io.debezium.connector.mongodb.transforms.outbox.MongoEventRouter

自定义配置

连接器可能会发出多种类型的事件消息(例如心跳消息、墓碑消息或关于事务的元数据消息)。要仅对发源于 outbox 集合的事件应用该转换,请定义一条 SMT 谓词语句,以便选择性地仅对这些事件应用该转换(/book/reader/debezium/mongodb-outbox-event-router#options-for-applying-the-transformation-selectively)。

选择性应用转换的选项

除 Debezium 连接器在数据库发生变更时发出的变更事件消息之外,连接器还会发出其他类型的消息,包括心跳消息,以及关于模式变更和事务的元数据消息。由于这些其他消息的结构与该 SMT 设计处理的变更事件消息结构不同,最好配置连接器以选择性地应用该 SMT,从而仅处理预期的数据变更消息。你可以使用以下任一方法来配置连接器,以便选择性地应用该 SMT:

  • 为该转换配置 SMT 谓词(/book/reader/debezium/using-smt-predicates-to-selectively-apply-transformations#applying-transformations-selectively)。
  • 使用该 SMT 的 route.topic.regex 配置选项。

将 Avro 用作负载格式

MongoDB outbox 事件路由 SMT 支持任意负载格式。outbox 集合中 payload 字段的值会被透明地传递下去。除了使用 JSON 之外,另一种选择是使用 Avro。这对于消息格式治理以及确保 outbox 事件模式能够以向后兼容的方式演进非常有益。

源应用程序如何为 outbox 消息负载生成 Avro 格式的内容不在本文档的讨论范围之内。一种可行的方法是利用 KafkaAvroSerializer 类来序列化 GenericRecord 实例。为确保 Kafka 消息的值正好是 Avro 二进制数据,请对连接器应用以下配置:

transforms=outbox,...
transforms.outbox.type=io.debezium.connector.mongodb.transforms.outbox.MongoEventRouter
value.converter=io.debezium.converters.ByteArrayConverter

默认情况下,payload 字段的值(Avro 数据)就是唯一的消息值。将 ByteArrayConverter 配置为值转换器后,payload 字段的值会原封不动地作为 Kafka 消息值传出。

请注意,这与为其他 SMT 建议使用的 BinaryDataConverter 不同,原因在于 MongoDB 在内部存储字节数组时采用了不同的方式。

Debezium 连接器可以配置为发出心跳事件、事务元数据事件或模式更改事件(具体支持情况因连接器而异)。这些事件无法由 ByteArrayConverter 序列化,因此必须提供额外的配置,以便转换器知道如何序列化这些事件。例如,以下配置演示了如何使用不带模式的 Apache Kafka JsonConverter:

transforms=outbox,...
transforms.outbox.type=io.debezium.connector.mongodb.transforms.outbox.MongoEventRouter
value.converter=io.debezium.converters.ByteArrayConverter
value.converter.delegate.converter.type=org.apache.kafka.connect.json.JsonConverter
value.converter.delegate.converter.type.schemas.enable=false

委托的 Converter 实现由 delegate.converter.type 选项指定。如果转换器还需要额外的配置选项,也可以一并指定,例如像上文那样使用 schemas.enable=false 来禁用 schema。

发送带附加字段的消息

你的 outbox 集合中可能包含一些字段,你希望将其值添加到发送的 outbox 消息中。例如,考虑一个 outbox 集合,其 aggregatetype 字段的值为 purchase-order,另一个字段 eventType 的可能取值为 order-created 和 order-shipped。附加字段可以使用 field:placement:alias 语法添加。

placement 的允许取值如下:

  • header
  • envelope
  • partition

要将 eventType 字段的值发送到 outbox 消息的头部,请按如下方式配置 SMT:

transforms=outbox,...
transforms.outbox.type=io.debezium.transforms.outbox.EventRouter
transforms.outbox.collection.fields.additional.placement=eventType:header:type

最终,Kafka 消息中将出现一个头部:以 type 作为键,以 eventType 字段的值作为值。

要在外发消息信封中输出 eventType 字段的值,请按如下方式配置 SMT:

transforms=outbox,...
transforms.outbox.type=io.debezium.transforms.outbox.EventRouter
transforms.outbox.collection.fields.additional.placement=eventType:envelope:type

要控制出站消息产生的分区,请按如下方式配置 SMT:

transforms=outbox,...
transforms.outbox.type=io.debezium.transforms.outbox.EventRouter
transforms.outbox.collection.fields.additional.placement=partitionField:partition

注意,对于 partition 放置方式,添加别名不会产生任何效果。

将转义的 JSON 字符串扩展为 JSON

默认情况下,Debezium outbox 消息的 payload 表示为一个字符串。当该字符串的原始来源是 JSON 格式时,生成的 Kafka 消息会使用转义序列表示该字符串,如下例所示:

# Kafka Topic: outbox.event.order
# Kafka Message key: "1"
# Kafka Message Headers: "id=596e275826f08b2730779e1f"
# Kafka Message Timestamp: 1556890294484
{
  "{\"id\": {\"$oid\": \"da8d6de63b7745ff8f4457db\"}, \"lineItems\": [{\"id\": 1, \"item\": \"Debezium in Action\", \"status\": \"ENTERED\", \"quantity\": 2, \"totalPrice\": 39.98}, {\"id\": 2, \"item\": \"Debezium for Dummies\", \"status\": \"ENTERED\", \"quantity\": 1, \"totalPrice\": 29.99}], \"orderDate\": \"2019-01-31T12:13:01\", \"customerId\": 123}"
}

你可以配置 outbox 事件路由器来展开消息内容,将转义的 JSON 还原为其原始的未转义格式。在转换后的字符串中,配套 schema 由原始 JSON 文档推断得出。以下示例展示了结果 Kafka 消息中展开后的 JSON:

# Kafka Topic: outbox.event.order
# Kafka Message key: "1"
# Kafka Message Headers: "id=596e275826f08b2730779e1f"
# Kafka Message Timestamp: 1556890294484
{
  "id": "da8d6de63b7745ff8f4457db", "lineItems": [{"id": 1, "item": "Debezium in Action", "status": "ENTERED", "quantity": 2, "totalPrice": 39.98}, {"id": 2, "item": "Debezium for Dummies", "status": "ENTERED", "quantity": 1, "totalPrice": 29.99}], "orderDate": "2019-01-31T12:13:01", "customerId": 123
}

要在该转换中启用字符串转换,请将 collection.expand.json.payload 的值设置为 true,并使用 StringConverter,如下例所示:

transforms=outbox,...
transforms.outbox.type=io.debezium.connector.mongodb.transforms.outbox.MongoEventRouter
transforms.outbox.collection.expand.json.payload=true
value.converter=org.apache.kafka.connect.storage.StringConverter

配置选项

下表描述了可以为 outbox 事件路由器 SMT 指定的选项。表中 Group 列表示该配置选项在 Kafka 中的分类。

表 2. outbox 事件路由器 SMT 配置选项说明

选项 默认值 分组 说明

collection.op.invalid.behavior

warn

Collection

决定当 outbox 集合上发生更新操作时 SMT 的行为。可设置的值包括:

  • warn - SMT 记录一条警告日志,并继续处理 outbox 集合中的下一个文档。
  • error - SMT 记录一条错误日志,并继续处理 outbox 集合中的下一个文档。
  • fatal - SMT 记录一条错误日志,然后连接器停止处理。

outbox 集合中的所有变更预期都应为插入或删除操作。也就是说,outbox 集合充当队列的角色;不允许更新 outbox 集合中的文档。SMT 会自动过滤掉 outbox 集合上的删除操作(用于移除已处理的 outbox 事件)。

collection.field.event.id

_id

Collection

指定包含唯一事件 ID 的 outbox 集合字段。该 ID 将以 id 为键存储在发出事件的头部信息中。

collection.field.event.key

aggregateid

Collection

指定包含事件键的 outbox 集合字段。当该字段包含值时,SMT 会将该值用作发出的 outbox 消息的键。这对于在 Kafka 分区中保持正确的顺序非常重要。

collection.field.event.timestamp

Collection

默认情况下,发出的 outbox 消息中的时间戳是 Debezium 事件的时间戳。要在 outbox 消息中使用其他时间戳,请将此选项设置为 outbox 集合中包含你希望出现在发出的 outbox 消息中的时间戳的字段。

collection.field.event.payload

payload

Collection

指定包含事件负载的 outbox 集合字段。

collection.expand.json.payload

false

Collection

指定 SMT 是否对字符串负载执行 JSON 展开。如果没有找到内容,或出现解析错误,则内容保持"原样"不变。

有关更多详细信息,请参见展开转义的 JSON章节。

collection.fields.additional.placement

Collection, Envelope

指定要添加到发件箱消息标头或信封中的一个或多个发件箱集合字段。请指定一个以逗号分隔的配对列表。在每个配对中,指定字段的名称,以及你希望该值位于标头中还是信封中。配对中的两个值用冒号分隔,例如:

id:header,my-field:envelope

要为字段指定别名,请指定一个三元组,其中别名作为第三个值,例如:

id:header,my-field:envelope:my-alias

第二个值是放置位置,它必须始终为 header 或 envelope。

配置示例参见在 Debezium 发件箱消息中发出附加字段。

collection.field.additional.missing

error

Collection, Envelope

确定当 collection.fields.additional.placement 属性指定的字段在发件箱负载中未找到时 SMT 的行为。可设置的值有:

  • error — 如果配置的附加字段缺失,SMT 将抛出异常。
  • ignore — SMT 将静默跳过缺失的字段。

collection.field.event.schema.version

Collection, Schema

设置后,此值将用作架构版本,如 Kafka Connect Schema Javadoc 中所述。

route.by.field

aggregatetype

Router

指定 outbox 集合中某个字段的名称。默认情况下,该字段中指定的值会成为连接器向其发出 outbox 消息的主题名称的一部分。有关示例,请参阅预期的 outbox 集合说明。

route.topic.regex

(?<routedByValue>.*)

Router

指定 outbox SMT 在 RegexRouter 中应用于 outbox 集合文档的正则表达式。该正则表达式是 route.topic.replacement SMT 选项设置的一部分。

当此属性设置为默认值 Router 时,SMT 会将 route.topic.replacement 属性中设置的默认 ${routedByValue} 变量,替换为 route.by.field 属性所指定的值。

route.topic.replacement

outbox.event​.${routedByValue}

Router

指定连接器向其发出 outbox 消息的主题名称。默认主题名称以字符串 outbox.event. 为前缀,后接分配给 outbox 集合文档的 aggregatetype 字段的值。

例如,如果 aggregatetype 的值为 customers,则连接器会将 outbox 消息发送到主题 outbox.event.customers。

要更改目标主题的默认名称,请执行以下任一操作:

route.tombstone.on.empty.payload

false

Router

指示空负载或 null 负载是否会导致连接器发出墓碑事件。

tracing.span.context.field

tracingspancontext

Tracing

包含追踪 span 上下文的字段名称。

tracing.operation.name

debezium-read

追踪

表示 Debezium 处理 span 的操作名称。

tracing.with.context.field.only

false

追踪

设置为 true 时,只会对包含已序列化上下文字段的事件进行追踪。

分布式追踪

Outbox 事件路由 SMT 支持分布式追踪。更多详情请参阅追踪文档。

评论

登录后参与评论

正在加载评论…