Skip to main content

概述

Debezium MongoDB CDC Source Connector 从支持变更流的 MongoDB 部署读取文档变化,并将文档的新增、更新和删除事件发布到 Kafka Topic。Topic 名称由逻辑前缀、数据库和集合组成;初始快照会先发布匹配集合中的现有文档,随后继续发布变更流事件。快照记录使用 READ 操作,新增、更新和删除分别对应 CREATE、UPDATE 和 DELETE 操作;删除事件默认还会产生 tombstone 记录。 它适用于将 MongoDB 中的业务数据实时送入数据湖、搜索、分析或下游服务,也适用于已有数据初始化后持续同步的 CDC 链路。源端恢复依赖 Kafka Connect 持久化的 Offset 以及 MongoDB 仍保留对应的变更流历史;消费者应按稳定键或来源位置处理可能的重复事件。

前置条件

  • MongoDB 必须部署为支持 Change Streams 的副本集或其他受支持的集群部署;独立模式 MongoDB 不能作为该 Connector 的变更流来源。
  • 用于连接 MongoDB 的账号必须具备读取目标数据库和集合的权限;如果启用信号或增量快照,还必须具备读取信号集合的权限。写入信号由负责发出信号的外部客户端完成,该客户端需要相应的写入权限。
  • 使用 TLS 时,Kafka Connect Worker 必须能够读取相应的信任库、密钥库及其密码,并且 MongoDB 服务端证书的主机名必须能够通过校验。

授权许可

使用 Apache License 2.0。

快速开始

准备可访问的 Kafka Connect 集群、Kafka 集群和 MongoDB 变更流源端,确认 Worker 能访问 MongoDB 和 Connector 插件目录,并准备一个有权读取源端的账号。创建和管理 Connector 时,参见 AutoMQ 的管理 Connector
<database-host><database-port><replica-set-name><mongodb-user><mongodb-password><logical-topic-prefix> 替换为实际值。默认 snapshot.mode=initial 时,如果没有可用的既有 Offset,Connector 会先对匹配到的源数据执行初始快照,再继续读取变更流;默认的 MongoDB 连接认证源是 admin。Converter 的具体格式应与消费者约定一致。

配置

连接与认证

topic.prefix

为 Connector 生成的 MongoDB Topic 提供唯一的逻辑前缀。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:高
  • 必填:是
  • 有效值 / 注意事项:只能包含字母、数字、连字符、点和下划线;同一 Kafka 集群中的不同 MongoDB Source 实例应使用不同前缀。

mongodb.connection.string

MongoDB 连接字符串。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:高
  • 必填:是
  • 有效值 / 注意事项:必须是有效的 MongoDB ConnectionString;Connector 校验配置时会建立实时连接进行验证。副本集连接应在连接字符串中提供副本集信息。

mongodb.user

MongoDB 用户名。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:高
  • 必填:否
  • 有效值 / 注意事项:当凭证已嵌入连接字符串,或使用自定义认证提供方时可以省略。

mongodb.password

MongoDB 密码。
  • 类型PASSWORD
  • 默认值:无固定默认值
  • 重要级别:高
  • 必填:否
  • 有效值 / 注意事项:可与 mongodb.user 配合使用;不要将真实密码写入日志、文档示例或版本库。

mongodb.authsource

保存 MongoDB 用户凭证的认证数据库。
  • 类型STRING
  • 默认值admin
  • 重要级别:中
  • 有效值 / 注意事项:当用户不在 admin 数据库中创建时,设置为实际认证数据库。

mongodb.ssl.enabled

是否为 MongoDB 连接启用 TLS/SSL。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:启用后应同时准备服务端证书校验所需的信任配置。

mongodb.ssl.invalid.hostname.allowed

是否允许 MongoDB 服务端证书的主机名校验失败。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项true 会关闭主机名保护并带来中间人攻击风险,仅适用于受控诊断环境。

mongodb.ssl.keystore

用于双向 TLS 认证的客户端密钥库路径。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:中
  • 必填:否
  • 有效值 / 注意事项:仅在 MongoDB 要求客户端证书时设置。

mongodb.ssl.keystore.password

客户端密钥库密码。
  • 类型PASSWORD
  • 默认值:无固定默认值
  • 重要级别:中
  • 必填:否
  • 有效值 / 注意事项:密钥库受密码保护时必须设置。

mongodb.ssl.keystore.type

客户端密钥库类型。
  • 类型STRING
  • 默认值PKCS12
  • 重要级别:中
  • 有效值 / 注意事项:仅在设置 mongodb.ssl.keystore 时生效,并应与密钥库实际格式一致。

mongodb.ssl.truststore

用于验证 MongoDB 服务端证书的信任库路径。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:中
  • 必填:否
  • 有效值 / 注意事项:在 Worker 不使用默认信任库,或需要专用 CA 时设置。

mongodb.ssl.truststore.password

信任库密码。
  • 类型PASSWORD
  • 默认值:无固定默认值
  • 重要级别:中
  • 必填:否
  • 有效值 / 注意事项:信任库需要解锁或验证密码时设置。

mongodb.ssl.truststore.type

信任库类型。
  • 类型STRING
  • 默认值PKCS12
  • 重要级别:中
  • 有效值 / 注意事项:仅在设置 mongodb.ssl.truststore 时生效,并应与信任库实际格式一致。

mongodb.connect.timeout.ms

建立 MongoDB 连接的超时时间,单位为毫秒。
  • 类型INT
  • 默认值10000
  • 重要级别:低
  • 有效值 / 注意事项:应根据网络延迟和部署位置调整;过小可能导致启动或重连失败。

mongodb.server.selection.timeout.ms

MongoDB 驱动选择可用服务端的超时时间,单位为毫秒。
  • 类型INT
  • 默认值30000
  • 重要级别:低
  • 有效值 / 注意事项:服务端故障转移或跨网络部署时可适当增大。

mongodb.socket.timeout.ms

MongoDB Socket 操作超时时间,单位为毫秒。
  • 类型INT
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项0 表示不设置 Socket 超时。

mongodb.heartbeat.frequency.ms

MongoDB 驱动集群监控心跳频率,单位为毫秒。
  • 类型INT
  • 默认值10000
  • 重要级别:低
  • 有效值 / 注意事项:影响驱动发现集群状态变化的节奏。

mongodb.poll.interval.ms

轮询 MongoDB 副本集成员变化的间隔,单位为毫秒。
  • 类型LONG
  • 默认值30000
  • 重要级别:中
  • 有效值 / 注意事项:必须为正数;需要更快发现拓扑变化时可缩短,但会增加轮询开销。

connection.validation.timeout.ms

Connector 配置校验阶段验证 MongoDB 连接的超时时间,单位为毫秒。
  • 类型LONG
  • 默认值60000
  • 重要级别:低
  • 有效值 / 注意事项:必须为正数;它只控制连接验证,不替代运行阶段的连接超时配置。

源范围与字段

database.include.list

指定要捕获的数据库范围。
  • 类型LIST
  • 默认值:无固定默认值
  • 重要级别:高
  • 有效值 / 注意事项:默认按正则表达式匹配;不能与 database.exclude.list 同时设置。使用名称匹配模式时按字面值匹配。

database.exclude.list

指定不捕获的数据库范围。
  • 类型LIST
  • 默认值:无固定默认值
  • 重要级别:高
  • 有效值 / 注意事项:默认按正则表达式匹配;不能与 database.include.list 同时设置。

collection.include.list

指定要捕获的集合范围。
  • 类型LIST
  • 默认值:无固定默认值
  • 重要级别:高
  • 有效值 / 注意事项:使用完整的 database.collection 名称或模式;不能与 collection.exclude.list 同时设置。

collection.exclude.list

指定不捕获的集合范围。
  • 类型LIST
  • 默认值:无固定默认值
  • 重要级别:高
  • 有效值 / 注意事项:使用完整的 database.collection 名称或模式;不能与 collection.include.list 同时设置。

field.exclude.list

从事件中排除指定字段。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:中
  • 有效值 / 注意事项:使用逗号分隔的 database.collection.field 或嵌套字段名称;数据库和集合支持 *

field.renames

重命名事件中的字段。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:中
  • 有效值 / 注意事项:使用逗号分隔的 旧字段:新字段 映射,并使用完整字段限定名;嵌套字段名称必须符合字段映射语法。

快照与初始化

snapshot.mode

选择初始快照和后续变更流的启动方式。
  • 类型STRING
  • 默认值initial
  • 重要级别:低
  • 有效值 / 注意事项:可选 alwaysinitialneverno_datainitial_onlywhen_neededconfiguration_basedcustomnever 已弃用,使用 no_data 替代。initial_only 完成快照后不会继续持续捕获变更;no_datanever 要求已有源状态及解码所需的结构前提。

snapshot.collection.filter.overrides

为指定集合声明快照读取过滤条件。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:中
  • 有效值 / 注意事项:使用逗号分隔的 database.collection 名称;每个名称都必须配套设置同名的 snapshot.collection.filter.overrides.<db>.<collection> 属性。

snapshot.delay.ms

开始快照前等待的时间,单位为毫秒。
  • 类型LONG
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负数;适合需要让源端或下游准备完成后再初始化的场景。

streaming.delay.ms

快照完成后开始读取变更流前等待的时间,单位为毫秒。
  • 类型LONG
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负数;只影响快照到持续捕获之间的等待。

snapshot.fetch.size

快照读取时使用的批量大小。
  • 类型INT
  • 默认值:无固定默认值
  • 重要级别:中
  • 有效值 / 注意事项:必须为非负整数;应结合文档大小、MongoDB 负载和 Worker 内存调整。

snapshot.max.threads

快照阶段用于并行读取集合的最大线程数。
  • 类型INT
  • 默认值1
  • 重要级别:中
  • 有效值 / 注意事项:必须为正数;实际线程数不超过集合数量。它增加的是单个 Connector 内部的快照并行度,不增加 Kafka Connect Task 数量。

snapshot.include.collection.list

指定需要执行快照的集合范围。
  • 类型LIST
  • 默认值:无固定默认值
  • 重要级别:中
  • 有效值 / 注意事项:使用正则表达式列表;它用于限制快照集合,不等同于持续变更流的捕获范围。

incremental.snapshot.chunk.size

增量快照每个分块读取的文档数量。
  • 类型INT
  • 默认值1024
  • 重要级别:中
  • 有效值 / 注意事项:必须为非负数;较小分块降低单次读取压力,但会增加增量快照管理开销。

incremental.snapshot.watermarking.strategy

选择增量快照的水位标记策略。
  • 类型STRING
  • 默认值INSERT_INSERT
  • 重要级别:低
  • 有效值 / 注意事项:可选 INSERT_INSERTINSERT_DELETE;必须启用信号能力才能使用信号驱动的增量快照。

缓冲与发布

max.batch.size

单次从内部队列取出并发布的最大记录数。
  • 类型INT
  • 默认值2048
  • 重要级别:中
  • 有效值 / 注意事项:必须为正数,且应小于 max.queue.size

max.queue.size

内部事件队列允许缓存的最大记录数。
  • 类型INT
  • 默认值8192
  • 重要级别:中
  • 有效值 / 注意事项:必须大于 max.batch.size;增大后可以吸收短时波动,但会增加内存占用和故障前缓存量。

max.queue.size.in.bytes

内部事件队列允许缓存的最大字节数。
  • 类型LONG
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:必须为非负数;0 表示不启用按字节限制。MongoDB 文档大小差异较大时可用它约束内存。

poll.interval.ms

Connector 从内部队列轮询记录的间隔,单位为毫秒。
  • 类型LONG
  • 默认值500
  • 重要级别:中
  • 有效值 / 注意事项:必须为正数;较小值通常降低发布等待,但会增加轮询频率。

query.fetch.size

快照或查询读取时使用的驱动抓取大小。
  • 类型INT
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:必须为非负整数;0 使用驱动或数据库默认行为。

事件与心跳

tombstones.on.delete

是否在删除事件后发送 tombstone 记录。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项true 发送 DELETE 和 tombstone;false 只发送 DELETE。下游压缩 Topic 或状态存储应按其删除约定选择。

skipped.operations

指定不发布的操作类型。
  • 类型LIST
  • 默认值[t]
  • 重要级别:低
  • 有效值 / 注意事项:可选操作码 cudtnone;默认跳过 truncate 操作。

heartbeat.interval.ms

发送心跳事件的间隔,单位为毫秒。
  • 类型INT
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:必须为非负数;0 禁用心跳事件。低流量源可启用心跳以提高进度和 Offset 可见性。

heartbeat.topics.prefix

心跳 Topic 的名称前缀。
  • 类型STRING
  • 默认值__debezium-heartbeat
  • 重要级别:低
  • 有效值 / 注意事项:只有启用心跳事件时才需要特别关注;应与 Topic 命名和权限策略一致。

extended.headers.enabled

是否在源记录中写入 Debezium 扩展上下文 Header。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:关闭后下游将无法依赖这些扩展 Header 获取相应上下文。

信号与增量快照

signal.data.collection

存放 Connector 信号的 MongoDB 数据库和集合。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:中
  • 有效值 / 注意事项:格式为 database.collection;未设置时禁用信号处理。使用增量快照前应创建该集合并授予相应权限。

signal.poll.interval.ms

轮询信号集合的间隔,单位为毫秒。
  • 类型LONG
  • 默认值5000
  • 重要级别:中
  • 有效值 / 注意事项:必须为正数;越小越快发现信号,但会增加源端读取频率。

signal.enabled.channels

启用的信号通道列表。
  • 类型LIST
  • 默认值[source]
  • 重要级别:中
  • 有效值 / 注意事项:默认启用 source 通道;只有配置了对应信号源和权限时才应增加其他通道。

Topic、Schema 与扩展

topic.naming.strategy

生成 Kafka Topic 名称的策略类。
  • 类型CLASS
  • 默认值io.debezium.schema.DefaultTopicNamingStrategy
  • 重要级别:中
  • 有效值 / 注意事项:类必须实现 TopicNamingStrategy;更换策略会改变 Topic 路由,可能影响下游订阅和历史数据兼容性。

schema.name.adjustment.mode

调整 Schema 名称以适配特定序列化格式的模式。
  • 类型STRING
  • 默认值none
  • 重要级别:低
  • 有效值 / 注意事项:可选 noneavroavro_unicode;选择前确认下游 Schema Registry 或 Converter 的命名要求。

sourceinfo.struct.maker

创建 MongoDB Source 信息结构的实现类。
  • 类型CLASS
  • 默认值io.debezium.connector.mongodb.MongoDbSourceInfoStructMaker
  • 重要级别:低
  • 有效值 / 注意事项:自定义类必须提供兼容的 SourceInfoStructMaker 实现;除非有明确的事件元数据扩展需求,否则保留默认值。

converters

注册自定义 Converter 及其配置前缀。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:低
  • 有效值 / 注意事项:使用逗号分隔的自定义 Converter 名称,并为每个 Converter 提供对应的类型和选项;自定义类必须在 Worker 插件路径中可见。

post.processors

注册事件后处理器及其配置前缀。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:低
  • 有效值 / 注意事项:使用逗号分隔的后处理器名称,并提供各自的类型和选项;后处理器会改变发布事件的内容或元数据。

错误处理

event.processing.failure.handling.mode

处理损坏或无法解析事件的方式。
  • 类型STRING
  • 默认值fail
  • 重要级别:中
  • 有效值 / 注意事项:可选 failwarnignorewarnignore 可能跳过无法处理的事件,应配合错误日志和业务完整性检查使用。

errors.max.retries

可重试源端错误的最大重试次数。
  • 类型INT
  • 默认值-1
  • 重要级别:低
  • 有效值 / 注意事项-1 表示不限次数,0 表示禁用重试,正数表示有限重试;重试不是端到端无重复交付保证。

errors.retry.timeout

Kafka Connect Worker 对失败记录或任务进行重试的总时长,单位为毫秒。
  • 类型LONG
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项0 表示禁用该重试窗口,-1 表示无限重试;它属于 Worker 错误策略,应与任务失败告警一起规划。

errors.tolerance

Kafka Connect Worker 对错误的容忍级别。
  • 类型STRING
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:可选 noneallall 会跳过问题记录,不能代替错误分析和补偿流程。

Kafka Connect 框架

connector.class

指定 Kafka Connect 要实例化的 Connector 类。
  • 类型STRING
  • 默认值:无固定默认值
  • 重要级别:高
  • 必填:是
  • 有效值 / 注意事项:使用 io.debezium.connector.mongodb.MongoDbConnector

tasks.max

Connector 可使用的最大 Kafka Connect Task 数量。
  • 类型INT
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:至少为 1;该 Connector 为 MongoDB 复制返回一个 Task 配置,增大此值不会为每个副本集创建独立的 Kafka Connect Task。副本集复制和快照并行度由 Connector 内部线程处理。

key.converter

序列化源记录 Key 的 Converter 类。
  • 类型CLASS
  • 默认值:无固定默认值
  • 重要级别:低
  • 有效值 / 注意事项:必须实现 Kafka Connect Converter;未在 Connector 级别设置时使用 Worker 默认值。

value.converter

序列化 MongoDB CDC 事件 Value 的 Converter 类。
  • 类型CLASS
  • 默认值:无固定默认值
  • 重要级别:高
  • 有效值 / 注意事项:必须实现 Kafka Connect Converter;应与下游消费者的反序列化方式一致。

最佳实践

首次接入时限定初始化范围并继续增量同步

适用业务场景:首次接入已有 MongoDB 数据,需要先建立指定数据库和集合的存量基线,再持续接收后续变化,避免把无关集合的存量数据一并发布。 配置示例:
关键说明:先确认集合范围、快照期间的源端负载和下游容量,再启动 Connector。在没有可用的既有 Offset 时,默认 snapshot.mode=initial 会发布匹配集合的 READ 记录并继续捕获变更;后续新增集合不会自动补发其已有存量,若需要补发应使用单独的快照规划。collection.include.listcollection.exclude.list 不能同时配置。

已有数据基线时只接入后续变化

适用业务场景:下游已经通过备份、批处理或其他方式建立了数据基线,只需要从 Connector 启动位置开始接收后续 MongoDB 变化,不希望再次发布现有文档。 配置示例:
关键说明:no_data 不发送现有文档的快照记录,但仍要求下游具备解析后续事件所需的结构和基线。切换到该模式前应明确下游基线对应的时间点,并保留稳定的 topic.prefix 和 Worker Offset;不要用它绕过未完成的初始快照或补回已经过期的 MongoDB 变更历史。

低流量源持续运行时提高进度可见性

适用业务场景:业务集合平时写入较少,但仍需要确认 Connector 持续读取、维护源端进度,并在下游监控中区分“没有变化”和“没有运行”。 配置示例:
关键说明:心跳事件用于提高低流量期间的进度和 Offset 可见性,不代表源端产生了业务变更。应同时监控 Task 状态、心跳 Topic 和 Worker Offset;重启后仍需依赖持久化 Offset 以及 MongoDB 变更流历史,心跳不会提供端到端 exactly-once 保证。

监控

监控内容

监控 Kafka Connect Worker、Connector 和 Task 的健康状态及重启次数,检查源端到 Kafka 的吞吐、端到端延迟、积压、Offset 提交时间和提交失败;同时关注 MongoDB 连接错误、重试、任务错误日志、事件处理失败、Kafka 发布失败以及 Worker JVM 的堆内存、GC、线程和 CPU。启用 errors.tolerance=all 或错误处理相关 DLQ 配置时,还应监控被跳过记录和 DLQ 流量,避免把错误容忍误判为正常同步。

导入 Grafana 大盘

下载 AutoMQ Connect Grafana Dashboard,在 Grafana 中选择已经采集 Kafka Connect 指标的数据源,确认集群、Worker、Connector 和 Task 标签与指标采集配置一致,然后通过 Dashboards 的导入功能上传 JSON 并映射数据源。

限制条件

  • 独立模式 MongoDB 不提供该 Connector 所需的 Change Streams,因此不能作为本 Connector 的变更流来源。
  • tasks.max 增大不会按副本集拆分为多个 Kafka Connect Task;副本集复制使用 Connector 内部线程,快照线程也不会改变 Task 数量。
  • initial_only 只执行快照,不会在快照结束后继续持续捕获变更。
  • no_data 和已弃用的 never 不会发布现有文档快照,并要求已有源状态及下游解码所需的结构前提。
  • 增量快照的窗口去重用于处理快照与并发变化的碰撞,不构成端到端 exactly-once 保证,也不代表单一全局时间点快照。
  • 变更流管道不能使用 MongoDB 索引;使用正则范围过滤时还会增加内部字段变换开销。
  • 连接重启或 Offset 恢复时,如果 MongoDB 已清理保存的变更流历史,Connector 不能从旧位置无间隙地重放所有中间事件。
  • 内部缓冲区不是持久化检查点;增大队列只会改变内存、背压和延迟行为,不改变交付保证。
  • 在 Kafka 发布完成但 Worker Offset 尚未持久化的故障窗口内,重启后可能重新发布已经发送过的事件;下游应能处理重复记录。

常见问题

为什么 Connector 连接 MongoDB 失败?

检查 MongoDB 是否为支持 Change Streams 的副本集或集群部署,连接字符串是否包含正确的主机、端口和副本集信息,认证数据库是否由 mongodb.authsource 正确指定,以及账号是否具备目标数据库读取权限。启用 TLS 时,再检查信任库、证书主机名和 Worker 文件权限;不要通过 mongodb.ssl.invalid.hostname.allowed=true 绕过生产环境证书校验。

为什么启动后没有收到现有文档?

确认 snapshot.mode 没有设置为 no_datanever 或其他不执行数据快照的模式,并检查 database.include.listcollection.include.list、排除列表和 snapshot.include.collection.list 是否排除了目标集合。若下游已经保存过同一 topic.prefix 的 Offset,重新启动不会因为重启自动重新发送完整存量;需要按数据基线和快照策略重新规划。

为什么只收到部分集合的事件?

检查数据库和集合包含列表、排除列表及其正则表达式,确认没有同时配置同一层级的 include 和 exclude。还要检查 snapshot.include.collection.list 是否只限制了快照,避免把快照范围和持续变更流范围混为一谈。

为什么删除事件后又出现一条空值记录?

这是 tombstones.on.delete=true 的默认行为:Connector 先发布 DELETE,再发布同一 Key 的 tombstone,便于 Kafka 日志压缩清理旧值。如果下游不接受 tombstone,可设置 tombstones.on.delete=false,但必须同步评估下游删除和状态清理逻辑。

为什么重启后出现重复事件?

Kafka Connect Offset 的持久化可能晚于事件发布;如果故障发生在发布和 Offset 刷新之间,恢复时可能从较早位置重新发布事件。确认 topic.prefix、Worker Offset 存储和 MongoDB 变更历史保持不变,并让下游使用稳定文档 Key、来源位置或业务幂等逻辑处理重复。Connector 不承诺端到端 exactly-once。

为什么启用增量快照后没有执行?

确认已设置 signal.data.collection,该集合存在且账号具备读取信号的权限,signal.enabled.channels 包含 source,并且增量快照配置使用了支持的水位策略。还要确认信号文档使用正确的数据库和集合名称,以及相关集合具备 Connector 能读取的文档键。