汇聚连接器

MongoDB 接收器

师成师成· 更新于 2026-09-28· 阅读 16 分钟· 0 次阅读

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

Debezium MongoDB Sink 连接器

概述

Debezium MongoDB sink 连接器从 Apache Kafka 主题中捕获变更事件记录,然后将这些记录转换为 MongoDB 文档,并写入指定 MongoDB sink 数据库中的集合。对于需要高可扩展性和快速数据检索的应用,将变更数据传播到基于集群的 MongoDB 环境(该环境使用分片和副本集等特性来优化读操作)可以显著提升检索性能。该连接器只能处理来自 Debezium 关系数据库连接器的变更事件。

有关与此连接器兼容的 MongoDB 版本信息,请参阅 Debezium 发布概览。

架构与工作原理

Debezium MongoDB sink 连接器通过消费 Kafka 主题中的 Debezium 变更事件,并将其作为文档级操作应用,从而将关系型源数据库中表的状态镜像到 MongoDB 数据库的集合中。

该连接器订阅的 Kafka 主题中填充了由 Debezium 关系数据库源连接器产生的事件消息。每条事件消息以结构化格式描述一次数据库操作(插入、更新或删除),其中包含事件的详细信息。连接器将传入的变更事件记录转换为 MongoDB 文档格式,然后将生成的文档写入目标 MongoDB 集合。

在接收到事件后,连接器会解析事件负载,并确定应将其发送到哪个 MongoDB 集合。根据事件负载中指定的事件类型,连接器随后会在目标集合中执行以下操作之一:

负载中的事件类型对应的操作
INSERT创建文档
UPDATE修改具有指定标识符的文档。
DELETE删除具有指定标识符的文档。

连接器使用 MongoDB Java 驱动程序与 MongoDB 数据库进行交互。

主题与 MongoDB 集合之间的映射关系由连接器配置决定。文档键作为文档的唯一标识符,确保更新、插入和删除操作被传播到正确的 MongoDB 文档和集合中,并且操作按正确的顺序执行。

通过将事件消息映射为 MongoDB 文档的这一过程,连接器能够将关系型数据库中表的状态镜像到 MongoDB 数据库的集合中。

限制

Debezium MongoDB sink 连接器具有以下限制:

仅支持关系型数据库 / RDBMS 源连接器

MongoDB sink 连接器只能消费由以下关系型数据库的 Debezium 连接器产生的变更事件:

  • MariaDB
  • MySQL
  • Oracle
  • PostgreSQL
  • SQL Server

该连接器无法处理来自任何其他 Debezium 连接器的变更事件消息,包括 Debezium MongoDB 源连接器。

Schema 演进

虽然该连接器能够处理基本的 schema 变更,但高级的 schema 演进场景可能需要人工干预或特定的配置。由于 MongoDB 是无 schema 的,其处理 schema 演进的能力有限。

事务支持

该连接器按照源系统中操作提交的先后顺序,以时间顺序逐条处理变更事件。虽然 MongoDB 支持事务,但 Debezium MongoDB 连接器并不提供跨多个 CDC 事件、或跨单个 sink 任务中多个文档的事务保证。

快速开始(使用 Kafka Connect)

部署一个基础的 MongoDB sink 连接器实例用于测试。

前提条件

以下组件在你的环境中可用且正在运行:

  • Kafka 集群
  • Kafka Connect
  • 一个 MongoDB 实例
  • Debezium 关系型数据库连接器
  • Debezium MongoDB 连接器

操作步骤

  1. 配置并启动一个 Debezium 源连接器(例如 Debezium PostgreSQL Connector),将关系型数据库的变更流式传输到 Kafka。

  2. 配置并启动 Debezium MongoDB sink 连接器,以消费源连接器发送到 Kafka 的事件,并将这些事件发送到 MongoDB sink 数据库。

    以下示例提供了 Debezium MongoDB sink 连接器的最小配置。请将示例中的占位符替换为你环境中实际的值。

    {
      "name": "mongodb-sink-connector",
      "config": {
        "connector.class": "io.debezium.connector.mongodb.sink.MongoDbSinkConnector",
        "topics.regex": "server1\.inventory\..*",
        "mongodb.connection.string": "mongodb://localhost:27017",
        "sink.database": "debezium"
      }
    }

配置

MongoDB sink 连接器接受多种配置选项,如下表所示。

属性默认值说明
connector.class无默认值必须设置为 io.debezium.connector.mongodb.sink.MongoDbSinkConnector。
tasks.max1任务的最大数量。
topics 或 topics.regex无默认值要消费的 Kafka 主题列表。如果你将此值设置为 topics.regex,连接器将消费所有与该正则表达式匹配的主题。

表 1. Kafka Connect sink 连接器必需的配置属性

表 2. MongoDB 连接必需的属性

属性 默认值 说明

mongodb.connection.string

无默认值

sink 用于连接 MongoDB 的 MongoDB 连接字符串(URI)。该 URI 遵循标准的 MongoDB 连接字符串格式。

示例:mongodb://localhost:27017/?replicaSet=my-replica-set

sink.database

无默认值

目标 MongoDB 数据库的名称。

表 3. sink 行为配置

属性 默认值 说明

collection.naming.strategy

io.debezium.sink.naming.DefaultCollectionNamingStrategy

指定连接器用于根据 Kafka 主题名称推导目标 MongoDB 集合名称的策略。

请指定以下值之一:

io.debezium.sink.naming.DefaultCollectionNamingStrategy

连接器直接从主题名称中获取表名,并将源主题名称中的点号(句点)字符替换为下划线。

自定义实现

你可以提供自己的 CollectionNameStrategy 实现。

collection.name.format

${topic}

根据 Kafka 主题名称推导目标集合名称的模板。

column.naming.strategy

io.debezium.sink.naming.DefaultColumnNamingStrategy

指定连接器用于为目标集合中的列命名的策略。

请指定以下值之一:

io.debezium.sink.naming.DefaultColumnNamingStrategy

使用原始字段名作为列名。

自定义实现

指定自定义的 CollectionNameStrategy 实现。

表 4. 常用 Sink 选项

属性 默认值 说明

field.include.list

空字符串

可选的、以逗号分隔的字段名列表,用于匹配变更事件值中要包含的字段的完全限定名称。字段的完全限定名称的形式为 fieldName 或 topicName:fieldName。

如果在配置中包含此属性,请不要设置 field.exclude.list 属性。

field.exclude.list

空字符串

可选的、以逗号分隔的字段名列表,用于匹配变更事件值中要排除的字段的完全限定名称。字段的完全限定名称的形式为 fieldName 或 topicName:fieldName。

如果在配置中包含此属性,请不要设置 field.include.list 属性。

batch.size

2048

单次批量写入的最大记录数。

配置示例

下面的示例演示了连接器的配置,该连接器从关系数据库中的特定主题读取变更事件,并将生成的文档写入 MongoDB Sink 数据库。示例中将连接器配置为捕获 dbserver1.inventory.customers、dbserver1.inventory.orders 和 dbserver1.inventory.products 主题中的变更事件,并将数据写入 MongoDB Sink 数据库中 debezium 集合的文档里。

下面的示例展示了如何配置连接器,使其从 dbserver1.inventory 数据库的三个特定主题读取变更事件,从而修改 MongoDB Sink 数据库中名为 debezium 的集合。

{
    "name": "mongodb-sink-connector",
    "config": {
        "connector.class": "io.debezium.connector.mongodb.sink.MongoDbSinkConnector",
        "topics": "dbserver1.inventory.customers,dbserver1.inventory.orders,dbserver1.inventory.products",
        "mongodb.connection.string": "mongodb://localhost:27017",
        "sink.database": "debezium"
    }
}

监控

此版本的连接器不提供任何指标。

键字段映射

当连接器处理事件时,它会将数据映射到目标 MongoDB 文档中的特定字段。

  • Debezium 变更事件中的键(例如 Kafka 消息键)默认映射到 MongoDB 的 _id 字段。
  • 值会被映射到 MongoDB 文档中。
  • 更新和删除操作基于键字段映射进行解析。

以下示例展示了 Kafka 主题中的事件键:

{
    "userId": 1,
    "orderId": 1
}

根据映射逻辑,前面的键值会被映射到 MongoDB 文档的 _id 字段,如下面的示例所示:

{
    "_id": {
        "userId": 1,
        "orderId": 1
    }
}

在 Debezium MongoDB Sink 连接器中使用 CloudEvents

Debezium MongoDB sink 连接器可以消费序列化为 CloudEvents 的记录。Debezium 能够以 CloudEvents 格式发出变更事件,使事件负载被封装在一个标准化的信封中。

当在源连接器上启用 CloudEvents 时,MongoDB sink 连接器会解析 CloudEvents 信封。

实际的 Debezium 事件负载会从 data 部分中提取出来。

随后,该事件会按照标准的插入、更新或删除语义应用到目标 MongoDB 集合中。

这一流程使得 Debezium 能够与更广泛的事件驱动系统集成,同时仍然将产生的事件持久化到 MongoDB 中。

属性默认值说明
cloud.events.schema.name.pattern.*CloudEvents\.Envelope$用于识别 CloudEvents 消息的正则表达式模式,通过该模式匹配架构名称。

表 5. CloudEvents sink 选项

评论

登录后参与评论

正在加载评论…