Skip to main content

概述

Couchbase Source Connector 从 Couchbase bucket 的 DCP 变更流读取文档 mutation、deletion 和过期事件,并将它们转换为 Kafka source record。一个 Connector 实例连接一个 bucket,可以读取全部 scope 和 collection,也可以限定到一个 scope 或一组 scope.collection,再按配置把记录路由到 Kafka Topic。 默认处理器输出带 Connect schema 的记录;也可以选择 Raw JSON 处理器,让 JSON mutation 以原始字节写入 Kafka,或选择带元数据的 Raw JSON 处理器。记录可以携带 document ID、bucket、scope、collection、vBucket、序列号、CAS 和过期时间等 Header,便于下游建立幂等处理和审计索引。 连接器保存每个 vBucket 的 source offset,并在重启或任务重新分配后从保存位置恢复。普通单记录路径应按至少一次交付设计:Kafka 写入后、offset 持久化前发生故障时,较早记录可能再次出现;顺序边界是同一 source vBucket,不保证跨 vBucket、Task 或 Topic 的全局顺序。

前置条件

  • Couchbase 账号必须能够访问目标 bucket,并具备读取 DCP 变更流所需的权限;使用 scope 或 collection 选择时,Couchbase Server 必须支持集合流。
  • 如果启用 TLS,Worker 必须能读取配置的 PEM CA 文件或 Java keystore;如果使用客户端证书认证,还要准备可读的客户端证书文件及其密码。
  • 如果使用 couchbase.persistence.polling.interval 的默认持久化检查,源 bucket 必须支持持久化;ephemeral bucket 必须把该配置设为 0
  • Connector 输出的目标 Topic、black-hole Topic、initial-offset Topic 或 schema 失败 Topic 必须预先准备并允许 Connect Worker 写入;连接器不会替你创建这些 Topic。
  • 选择 RawJsonSourceHandlerRawJsonWithMetadataSourceHandler 时,Worker 或 Connector 必须使用与原始字节输出匹配的 ByteArrayConverter

授权许可

使用 Apache License 2.0。

快速开始

准备 Connect Cluster、Kafka 集群和可访问的 Couchbase bucket,确认 Worker 到 Couchbase、Kafka 及目标 Topic 的网络连通和访问权限。以下是 Connector properties 配置;把尖括号中的地址、资源名和凭据替换为实际值。创建和管理 Connector 的通用步骤请参阅管理 Connector
couchbase.seed.nodes 设置为一个或多个 Couchbase KV 节点地址,couchbase.bucket 设置为源 bucket,couchbase.topic 设置为已准备好的 Kafka Topic。该示例发布 JSON mutation 的原始字节;删除和过期事件的 value 为 null。生产环境不要把真实密码写入文档、日志或共享配置,应用时应使用安全的配置注入方式。

配置

连接与身份验证

couchbase.seed.nodes

Couchbase 集群的种子节点地址列表。
  • 类型list
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必填。使用逗号分隔的节点地址;自定义端口应为 KV 端口,通常为非 TLS 的 11210 或 TLS 的 11207
  • 必填:是

couchbase.username

连接 Couchbase 使用的用户名。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必填。使用客户端证书认证时,运行时可忽略该值,但仍需提供该配置项以满足配置定义。
  • 必填:是

couchbase.password

连接 Couchbase 使用的密码。
  • 类型password
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必填。可通过 KAFKA_COUCHBASE_PASSWORD 在读取配置时提供;使用客户端证书认证时运行时可忽略该值,但仍需提供配置项。
  • 必填:是

couchbase.bucket

要读取变更流的 Couchbase bucket。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:运行时必填,不能留空。一个 Connector 实例只连接一个 bucket。
  • 必填:是

couchbase.network

选择 Couchbase SDK 的网络解析方式。
  • 类型string
  • 默认值auto
  • 重要级别:中
  • 有效值 / 注意事项:区分大小写,可选 autodefaultexternal

couchbase.bootstrap.timeout

连接器启动和 DCP socket 建立使用的超时时间。
  • 类型string
  • 默认值30s
  • 重要级别:中
  • 有效值 / 注意事项:使用整数加 mssmhd,例如 10s

TLS 与客户端证书

couchbase.enable.tls

是否启用 Couchbase TLS 连接。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:启用后会使用信任证书、主机名校验和可选客户端证书配置;未指定自定义信任材料时使用默认 CA。

couchbase.enable.hostname.verification

是否校验 TLS 证书中的主机名。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项:仅对 TLS 连接生效。除非有明确的证书部署原因,不要关闭主机名校验。

couchbase.trust.store.path

Java keystore 信任库的绝对路径。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:与 couchbase.trust.store.password 配套使用;也可以改用 couchbase.trust.certificate.path

couchbase.trust.store.password

Java keystore 信任库密码。
  • 类型password
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:配置 couchbase.trust.store.path 时提供。可通过 KAFKA_COUCHBASE_TRUST_STORE_PASSWORD 覆盖。

couchbase.trust.certificate.path

PEM 格式 CA 证书的绝对路径。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:与 trust store 二选一;Worker 进程必须能读取该文件。

couchbase.client.certificate.path

客户端证书 keystore 或 PKCS12 bundle 的路径。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:非空时启用客户端证书认证,并忽略运行时的用户名和密码值。

couchbase.client.certificate.password

客户端证书文件密码。
  • 类型password
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:与 couchbase.client.certificate.path 配套使用。可通过 KAFKA_COUCHBASE_CLIENT_CERTIFICATE_PASSWORD 覆盖。

Topic 路由与记录处理

couchbase.topic

未配置集合级覆盖时使用的目标 Kafka Topic 模板。
  • 类型string
  • 默认值${bucket}.${scope}.${collection}
  • 重要级别:中
  • 有效值 / 注意事项:支持 ${bucket}${scope}${collection} 占位符;也可使用 couchbase.topic[scope.collection] 为某个集合覆盖。

couchbase.collection.to.topic

scope.collection=topic 形式配置集合到 Topic 的映射。
  • 类型list
  • 默认值:空列表
  • 重要级别:中
  • 有效值 / 注意事项:已弃用。新配置使用 couchbase.topic[scope.collection];未映射集合回退到 couchbase.topic
  • 已弃用:是
  • 替代项couchbase.topic[scope.collection]

couchbase.source.handler

把 Couchbase 文档事件转换为 Kafka source record 的处理器类。
  • 类型class
  • 默认值:无
  • 重要级别:中
  • 有效值 / 注意事项:必填。类必须实现 SourceHandlerMultiSourceHandler 并可实例化;使用 RawJsonSourceHandler 时,value.converter 应设为 org.apache.kafka.connect.converters.ByteArrayConverter
  • 必填:是

couchbase.headers

选择附加到 Kafka record 的 Couchbase 元数据 Header。
  • 类型list
  • 默认值:空列表
  • 重要级别:中
  • 有效值 / 注意事项:可选值为 bucketscopecollectionkeyqualifiedKeycaspartitionpartitionUuidseqnorevexpiry

couchbase.header.name.prefix

couchbase.headers 选中的 Header 名称添加前缀。
  • 类型string
  • 默认值couchbase.
  • 重要级别:中
  • 有效值 / 注意事项:前缀会添加到所有选中的 Header 名称前。

couchbase.event.filter

过滤 Couchbase 事件的 Filter 类。
  • 类型class
  • 默认值com.couchbase.connect.kafka.filter.AllPassFilter
  • 重要级别:中
  • 有效值 / 注意事项:类必须可实例化并实现 Filter。默认过滤器排除系统 scope 事件及 key 以 _txn: 开头的事务元数据文档。

couchbase.jsonpath.filter

按文档 JSON 内容筛选 mutation。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:空值表示不启用 JSONPath 筛选;非空值必须是可解析的 JSONPath。它在事件 Filter 和 source handler 之前执行,可使用 couchbase.jsonpath.filter[scope.collection] 覆盖。

couchbase.no.value

是否让 DCP 省略 mutation 的文档正文。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:启用后保留 key 和元数据但不提供正文;依赖正文的 JSONPath、Filter、Handler 和 schema 映射将无法生成正文内容。

couchbase.xattrs

是否请求 Couchbase extended attributes。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:启用后,Filter 或自定义 source handler 可以读取 xattrs;只有确实需要这些元数据时才启用。

起始位置与 Offset

couchbase.stream.from

设置首次启动或缺少保存 Offset 的 vBucket 的读取起点。
  • 类型string
  • 默认值SAVED_OFFSET_OR_BEGINNING
  • 重要级别:中
  • 有效值 / 注意事项:区分大小写,可选 SAVED_OFFSET_OR_BEGINNINGSAVED_OFFSET_OR_NOWBEGINNINGNOWBEGINNING 受 Couchbase 保留历史边界限制。

couchbase.initial.offset.topic

为没有保存 Offset 的 vBucket 发布初始位置合成记录的 Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:仅在 couchbase.stream.from=SAVED_OFFSET_OR_NOW 时相关;配置 couchbase.black.hole.topic 时应使用同一个 Topic。该 Topic 不是业务数据 Topic。

couchbase.black.hole.topic

接收被 Filter 或 Handler 忽略事件的合成记录,以便提交其 source offset。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:非空时为每个忽略事件发送小型占位记录;应按业务需要为该 Topic 配置较短保留时间和较小分段。它不是业务 DLQ。

couchbase.connector.name.in.offsets

是否把 Connector 名称写入 source offset 标识。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:已弃用。新部署保持 false;只有为了兼容曾经以 Connector 名称保存 Offset 的旧部署时才设为 true
  • 已弃用:是
  • 替代项:新部署保持 false,使用 Kafka Connect 默认的 Connector offset 隔离。

批处理与 DCP

couchbase.batch.size.max

单次 poll 最多返回的 SourceRecord 数量。
  • 类型int
  • 默认值2000
  • 重要级别:中
  • 有效值 / 注意事项:只限制单次 poll 批次,不限制 Task 内部事件队列或 Kafka producer 的在途数据;应结合 Worker 吞吐和堆内存调整。

couchbase.compression

选择 DCP 传输压缩方式。
  • 类型string
  • 默认值ENABLED
  • 重要级别:中
  • 有效值 / 注意事项:区分大小写,可选 DISABLEDFORCEDENABLED。只影响 Couchbase 到 Connector 的传输,不改变 Kafka record 压缩。

couchbase.persistence.polling.interval

等待 Couchbase 变更满足持久化条件后再发布的轮询间隔。
  • 类型string
  • 默认值100ms
  • 重要级别:中
  • 有效值 / 注意事项:使用整数加 mssmhd;设为 0 可关闭该检查。ephemeral bucket 必须设为 0,否则事件可能无法发布。

couchbase.flow.control.buffer

每个 Task 在每个 Couchbase 节点上使用的 DCP 流控缓冲区大小。
  • 类型string
  • 默认值16m
  • 重要级别:中
  • 有效值 / 注意事项:支持 bkmg;缓冲区会按 Task 和节点数量放大,增大它会提高内存占用。常见调优范围为 10m50m,仍需结合实际负载评估。

Schema 处理

couchbase.value.schema

供支持 schema 的 source handler 使用的 Avro record schema JSON。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:仅对支持该配置的 handler 生效;ConfigurableSchemaSourceHandler 要求可解析的 Avro record schema。可使用 couchbase.value.schema[scope.collection] 按集合覆盖。

couchbase.schema.failure.action

Schema Registry handler 遇到 schema 缺失或不匹配时的处理方式。
  • 类型string
  • 默认值TERMINATE
  • 重要级别:中
  • 有效值 / 注意事项:区分大小写,可选 TERMINATEDROPDLQ;只对支持该选项的 Schema Registry handler 生效。

couchbase.dlq.topic

Schema Registry handler 在 couchbase.schema.failure.action=DLQ 时使用的目标 Topic。
  • 类型string
  • 默认值couchbase.dlq
  • 重要级别:中
  • 有效值 / 注意事项:仅用于该 handler 的 schema 失败转发,不等同于 Kafka Connect 通用错误处理框架的 DLQ;目标 Topic 必须允许 Worker 写入。

日志与诊断

couchbase.log.redaction

控制 Couchbase 客户端日志中的敏感信息脱敏级别。
  • 类型string
  • 默认值NONE
  • 重要级别:中
  • 有效值 / 注意事项:区分大小写,可选 NONEPARTIALFULL;按组织的日志安全要求选择。

couchbase.log.document.lifecycle

是否把每个文档的生命周期里程碑从 DEBUG 提升到 INFO。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:启用后可能产生较高日志量,适合短时间排查单条文档流转。

couchbase.metrics.interval

周期性写入 Connector 日志的指标间隔。
  • 类型string
  • 默认值10m
  • 重要级别:中
  • 有效值 / 注意事项:使用整数加 mssmhd;设为 0 禁用指标日志。该配置属于未承诺稳定性的诊断项。

couchbase.enable.dcp.trace

是否启用详细 DCP trace 日志。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:启用后以 INFO 级别输出详细 DCP 诊断,并启用 couchbase.dcp.trace.document.id.regex;只建议在故障排查期间短时使用。

couchbase.dcp.trace.document.id.regex

DCP trace 要记录的文档 ID 正则表达式。
  • 类型string
  • 默认值.*
  • 重要级别:中
  • 有效值 / 注意事项:仅在 couchbase.enable.dcp.trace=true 时使用;连接器启动时编译 Java 正则表达式。

集合范围

couchbase.scope

选择一个 scope 下的全部 collection。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:需要 Couchbase Server 7.0 或更高版本。与 couchbase.collections 互斥;两者都为空时读取全部 scope 和 collection。

couchbase.collections

选择要读取的限定集合名称列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:中
  • 有效值 / 注意事项:使用逗号分隔的 scope.collection 名称,需要 Couchbase Server 7.0 或更高版本。与 couchbase.scope 互斥;两者都为空时读取全部 scope 和 collection。

动态集合上下文配置

couchbase.topic[scope.collection]

为指定 scope.collection 覆盖 couchbase.topic 模板。
  • 类型string
  • 默认值:继承 couchbase.topic${bucket}.${scope}.${collection}
  • 重要级别:中
  • 有效值 / 注意事项:只有提交配置中出现格式正确的方括号属性时注册;方括号内必须是限定的 scope 和 collection,未匹配的集合回退到基础配置。

couchbase.jsonpath.filter[scope.collection]

为指定 scope.collection 覆盖 JSONPath 筛选表达式。
  • 类型string
  • 默认值:继承 couchbase.jsonpath.filter 的空字符串
  • 重要级别:中
  • 有效值 / 注意事项:每个覆盖值都必须是可解析的 JSONPath;未匹配的集合使用基础配置。

couchbase.value.schema[scope.collection]

为指定 scope.collection 覆盖 handler 使用的 schema JSON。
  • 类型string
  • 默认值:继承 couchbase.value.schema 的空字符串
  • 重要级别:中
  • 有效值 / 注意事项:仅对支持该配置的 source handler 生效;每个覆盖值都必须符合该 handler 的 schema 要求。

Kafka Connect 框架

connector.class

选择要运行的 Kafka Connect Source 插件类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必填,使用 com.couchbase.connect.kafka.CouchbaseSourceConnector
  • 必填:是

tasks.max

允许该 Connector 创建的最大 Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:至少为 1。连接器按 Couchbase vBucket 切分任务,实际非空 Task 数不会超过 vBucket 数;同一 vBucket 不会由该 Connector 的多个 Task 并行消费。

key.converter

序列化 Source record key 的 Converter。
  • 类型class
  • 默认值null(显式值,继承 Worker 配置)
  • 重要级别:低
  • 有效值 / 注意事项:必须是可实例化的 Kafka Connect Converter;Source handler 通常使用文档 ID 作为 key。示例使用 org.apache.kafka.connect.storage.StringConverter

value.converter

序列化 Source record value 的 Converter。
  • 类型class
  • 默认值null(显式值,继承 Worker 配置)
  • 重要级别:低
  • 有效值 / 注意事项:必须是可实例化的 Kafka Connect Converter;Raw JSON handler 需要 org.apache.kafka.connect.converters.ByteArrayConverter,其他 handler 应按输出 schema 选择匹配的 Converter。

最佳实践

新建实时链路时只接收启用后的变更

适用于新建实时事件链路,但不需要回放 Couchbase 当前仍保留的历史变更的场景。目标是让已有 Offset 的 vBucket 正常恢复,让尚无 Offset 的 vBucket 从 Connector 启用时的位置开始,并把该初始位置纳入 Kafka Connect 的 Offset 提交链路。
关键说明:预先创建 <offset-topic> 并允许 Worker 写入。SAVED_OFFSET_OR_NOW 不会为没有保存 Offset 的 vBucket 回放既有历史;couchbase.initial.offset.topic 会写入初始位置合成记录,待该记录和 source offset 提交成功后形成恢复点。它不是业务数据 Topic,不应由业务消费者处理。

按集合限定源范围并路由到不同 Topic

适用于一个 bucket 中只有部分集合需要进入 Kafka,且不同业务集合需要隔离到不同 Topic 的场景。源端范围在 DCP 流建立时限定,Topic 使用集合上下文覆盖,避免把无关集合先读入再在 Kafka 侧过滤。
关键说明:couchbase.collections 使用完整的 scope.collection 名称;不要把它写成通配符,也不要同时配置互斥的 couchbase.scope。Topic 模板会把两个集合分别路由到 orders-pendingorders-completed;请预先创建对应 Topic。只有个别集合需要例外名称时,再使用配置章节说明的 couchbase.topic[scope.collection] 覆盖。带元数据的 Raw JSON 输出便于消费者同时使用原始正文和事件上下文。

在故障一致性与吞吐之间选择取舍

适用于持久化 bucket 需要降低源端 failover 回滚事件进入 Kafka 的概率,或 ephemeral bucket 需要立即发布事件的场景。根据 bucket 类型和可接受延迟选择持久化检查,配合 DCP 流控缓冲区控制吞吐和内存。
关键说明:持久化检查会增加延迟、网络和内存开销,但可降低源端 failover 后出现 alternate-history 事件的风险;设为 0 会优先吞吐和低延迟,但下游必须接受这种风险。ephemeral bucket 必须使用 0,否则持久化条件无法满足;缓冲区会按 Task 和 Couchbase 节点数量放大,不能仅按单个配置值估算堆内存。

监控

监控内容

监控 Kafka Connect Worker、Connector 和 Task 的健康状态与重启次数,观察 Source 吞吐、端到端延迟、poll 延迟、source offset 提交进度、错误和重试;同时关注 Worker JVM 堆使用、GC、线程和 CPU。若启用了错误处理或 Schema Registry handler 的专用失败 Topic,再监控对应 DLQ 或失败 Topic 的写入量和积压,并将其与目标 Topic 的生产速率、Consumer Lag 和 Couchbase 源端负载一起分析。

导入 Grafana 大盘

下载 AutoMQ Connect Cluster Dashboard,在 Grafana 中选择已接入 Connect Worker 指标的 Prometheus 数据源,并确认指标标签与大盘变量匹配,然后使用 Grafana 的导入功能加载 JSON。该大盘用于 Kafka Connect 集群级指标,Connector 特有诊断仍需结合 Worker 日志和 Couchbase 监控。

限制条件

  • 一个 Connector 实例只能连接一个 Couchbase bucket;跨 bucket 需要创建独立的 Connector 实例。
  • 连接器只保证同一 source vBucket 内的顺序,不保证跨 vBucket、Task、Topic 或 collection 的全局顺序。
  • 普通单记录路径应按至少一次交付处理;Kafka 写入后、source offset 持久化前的故障以及 Couchbase failover rewind 都可能造成重复记录。
  • Task 内部事件队列没有按元素数量设置的硬上限;couchbase.flow.control.buffer 是 DCP 字节流控,不是 JVM 队列容量上限。
  • 连接器没有统一覆盖连接、Filter、Handler、schema 和 producer 失败的次数或退避重试配置;恢复行为依赖 DCP 客户端和 Kafka Connect Task 生命周期。

常见问题

Connector 已启动但 Kafka 中没有业务消息,应该检查什么?

先检查 Couchbase 地址、bucket、账号权限和目标 Topic 写权限,再确认 couchbase.source.handler 可实例化且 couchbase.stream.from 与预期起点一致。若配置了 couchbase.scopecouchbase.collections 或 JSONPath,确认目标文档确实属于所选范围并满足表达式;默认 Filter 还会排除系统 scope 和 _txn: 事务元数据文档。最后检查 Task 日志中的 DCP 连接错误和 Kafka producer 错误。

为什么重启后会再次看到已经发布过的记录?

Source offset 在 Kafka record 成功写入后才由 Kafka Connect 持久化;如果 Worker 在写入后、Offset 提交前中断,恢复时会从较早位置重放。Couchbase failover 发生 rewind 时也可能重复。消费者应使用稳定的文档 ID、vBucket 与 seqno 等元数据实现幂等处理,并避免随意更改 Connector 身份或已保存 Offset 的兼容配置。

为什么配置了 couchbase.no.value=true 后正文为空?

该配置会让 DCP 省略 mutation 正文,只保留 key 和元数据,适合只需要事件标识或元数据的处理。需要按业务字段过滤、生成 JSON 或构造 schema record 时,应恢复为 false,并确认 couchbase.jsonpath.filter、Filter 和 source handler 不依赖被省略的正文。

为什么启用持久化轮询后事件延迟,甚至完全没有事件?

持久化轮询会等待源端变更满足持久化条件,因此会增加延迟和内存占用。持久化 bucket 可以根据一致性要求调整 couchbase.persistence.polling.interval;ephemeral bucket 没有持久化条件,必须把该值设为 0。同时检查 DCP 流控缓冲区、Task 数和 Couchbase 节点状态。

为什么不同 collection 的记录进入了同一个 Topic?

检查 couchbase.topic 模板是否包含 ${scope}${collection},以及集合级属性是否使用完整的 couchbase.topic[scope.collection] 格式。集合级覆盖只对方括号中的限定名称生效;未匹配的集合会回退到基础 Topic。连接器不会自动创建或校验 Topic 名称。

使用 Schema Registry handler 时,schema 失败会怎样处理?

检查 couchbase.schema.failure.actionTERMINATE 会让任务失败,DROP 会丢弃不匹配事件,DLQ 会把诊断记录发送到 couchbase.dlq.topic。这些配置只对支持它们的 Schema Registry handler 生效,且目标 Topic、value Converter 和 Schema Registry 访问必须匹配;它不等同于 Kafka Connect 通用 DLQ。