Skip to main content

概述

Couchbase Sink Connector 从 Kafka Topic 消费记录,并将记录写入 Couchbase。默认处理器把非空记录值转换为 JSON 文档并执行整文档 upsert;墓碑记录用于删除对应文档。Connector 也支持 N1QL、Analytics 和子文档处理器,以适应条件更新、分析数据写入和局部文档变更等场景。 Connector 可以按 Kafka Topic 将记录路由到不同的 Couchbase Bucket、Scope 和 Collection。文档 ID 可以来自消息值中的业务字段、Kafka 记录键或 Topic、分区与 Offset 的组合。目标写入完成后,Kafka Connect 才推进该批记录的 Offset,但 Couchbase 写入与 Kafka Offset 提交不构成跨系统事务。

前置条件

  • 对于默认 KV、N1QL 和子文档处理器,目标 Bucket、Scope 和 Collection 必须已存在,连接账号需要具备相应的读写权限;使用 Analytics 处理器时,还需要可用的 Couchbase Analytics 服务和目标数据集。
  • 使用 TLS 时,需要准备受信任的 PEM CA 证书或 Java 信任库;使用双向 TLS 时,还需要客户端证书包及其密码。
  • couchbase.seed.nodes 使用 Connect Worker 可访问的 Couchbase 节点地址;指定自定义端口时必须使用 KV 端口。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka、目标 Couchbase Bucket 和 Collection,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector
<topic-name><couchbase-host><username><password><bucket-name> 替换为实际值。默认处理器将非空记录写入目标 Bucket 的 _default._default Collection;使用其他 Scope 或 Collection 时,增加 couchbase.default.collection

配置

Kafka 输入与任务

connector.class

指定要加载的 Couchbase Sink Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用 com.couchbase.connect.kafka.CouchbaseSinkConnector
  • 必填:是

topics

指定 Connector 消费的 Kafka Topic 列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:高
  • 有效值 / 注意事项:使用逗号分隔 Topic;必须在 topicstopics.regex 中恰好配置一个非空值。

topics.regex

使用 Java 正则表达式选择 Connector 消费的 Kafka Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:填写有效的 Java Pattern;必须在 topicstopics.regex 中恰好配置一个非空值。

tasks.max

指定 Kafka Connect 可为该 Connector 启动的最大 Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:必须不小于 1;有效并行度还受输入分区数量和 Couchbase 处理能力限制。

key.converter

指定 Kafka 记录键的 Connector 级 Converter。
  • 类型class
  • 默认值null,继承 Worker 配置
  • 重要级别:低
  • 有效值 / 注意事项:使用可实例化的 Kafka Connect Converter;记录键可作为 Couchbase 文档 ID 的回退来源。

value.converter

指定 Kafka 记录值的 Connector 级 Converter。
  • 类型class
  • 默认值null,继承 Worker 配置
  • 重要级别:低
  • 有效值 / 注意事项:使用与输入数据兼容的 Converter;转换结果必须能被 Connector 序列化为 JSON。

Couchbase 连接与认证

couchbase.seed.nodes

指定用于发现 Couchbase 集群的引导节点。
  • 类型list
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:填写一个或多个可访问的节点地址并以逗号分隔;自定义端口必须是 KV 端口。
  • 必填:是

couchbase.username

指定密码认证使用的 Couchbase 用户名。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:配置解析要求提供该值;使用客户端证书认证时,运行时会忽略该值。
  • 必填:是

couchbase.password

指定密码认证使用的 Couchbase 密码。
  • 类型password
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:配置解析要求提供该值;环境变量 KAFKA_COUCHBASE_PASSWORD 可在解析后覆盖它。使用客户端证书认证时,运行时会忽略该值。
  • 必填:是

couchbase.bucket

指定默认目标 Bucket,并作为省略 Bucket 的 Collection Keyspace 的默认 Bucket。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:使用 KV Collection 的处理器必须配置已存在且可访问的 Bucket;Analytics 处理器可以不配置。

couchbase.network

选择 Couchbase SDK 使用的集群地址网络视图。
  • 类型string
  • 默认值auto
  • 重要级别:中
  • 有效值 / 注意事项:可选 autodefaultexternal

couchbase.bootstrap.timeout

设置启动时建立 Couchbase 连接的最长等待时间。
  • 类型string
  • 默认值30s
  • 重要级别:中
  • 有效值 / 注意事项:使用非负整数加 mssmhd,也可使用 0

TLS 与证书

couchbase.enable.tls

控制是否为 Couchbase 连接启用 TLS。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:可选 truefalse;安全端点应启用 TLS 并配置适用的信任材料。

couchbase.enable.hostname.verification

控制是否校验 TLS 服务端证书中的主机名。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项:可选 truefalse;关闭会削弱服务端身份校验。

couchbase.trust.store.path

指定包含受信任 Couchbase CA 的 Java KeyStore 绝对路径。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:路径必须可读;它是 couchbase.trust.certificate.path 的信任材料替代方案。

couchbase.trust.store.password

指定信任库密码。
  • 类型password
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:按信任库要求填写;环境变量 KAFKA_COUCHBASE_TRUST_STORE_PASSWORD 可覆盖该值。

couchbase.trust.certificate.path

指定受信任 Couchbase CA 的 PEM 证书绝对路径。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:路径必须可读;它是 couchbase.trust.store.path 的 PEM 替代方案。

couchbase.client.certificate.path

指定包含客户端私钥和证书链的 Java KeyStore 或 PKCS12 文件路径。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:设置非空值会选择客户端证书认证,并在运行时忽略用户名和密码。

couchbase.client.certificate.password

指定客户端证书包密码。
  • 类型password
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:按证书包要求填写;环境变量 KAFKA_COUCHBASE_CLIENT_CERTIFICATE_PASSWORD 可覆盖该值。

路由与文档标识

couchbase.default.collection

指定默认目标 Keyspace,并允许按 Kafka Topic 覆盖。
  • 类型string
  • 默认值_default._default
  • 重要级别:中
  • 有效值 / 注意事项:使用 scope.collectionbucket.scope.collection;省略 Bucket 时使用 couchbase.bucket。按 Topic 覆盖使用 couchbase.default.collection[<kafka-topic>],未覆盖的 Topic 继承基础值。

couchbase.topic.to.collection

使用旧式映射按 Topic 指定目标 Keyspace。
  • 类型list
  • 默认值:空列表
  • 重要级别:中
  • 有效值 / 注意事项:条目格式为 topic=scope.collectiontopic=bucket.scope.collection;优先使用 couchbase.default.collection[<kafka-topic>]
  • 弃用状态:已弃用
  • 替代项couchbase.default.collection[<kafka-topic>]

couchbase.document.id

设置从记录值提取 Couchbase 文档 ID 的格式。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:可使用 ${/id}prefix::${/id} 形式的 JSON Pointer 占位符。空值或提取失败时回退到 Kafka 记录键,再回退到 Topic、分区和 Offset。按 Topic 覆盖使用 couchbase.document.id[<kafka-topic>]

couchbase.topic.to.document.id

使用旧式映射按 Topic 指定文档 ID 格式。
  • 类型list
  • 默认值:空列表
  • 重要级别:中
  • 有效值 / 注意事项:条目格式为 topic=<document-id-format>,每个条目只能包含一个映射分隔符;优先使用 couchbase.document.id[<kafka-topic>]
  • 弃用状态:已弃用
  • 替代项couchbase.document.id[<kafka-topic>]

couchbase.remove.document.id

控制是否从写入的文档中移除用于生成文档 ID 的字段。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅在 couchbase.document.id 的所有 JSON Pointer 都成功提取时有意义。按 Topic 覆盖使用 couchbase.remove.document.id[<kafka-topic>]

couchbase.document.expiration

设置文档过期时间。
  • 类型string
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:使用非负整数加 mssmhd0 表示不过期。按 Topic 覆盖使用 couchbase.document.expiration[<kafka-topic>]

写入处理与重试

couchbase.sink.handler

指定将 Sink 记录转换为 Couchbase 操作的 SinkHandler 实现。
  • 类型class
  • 默认值com.couchbase.connect.kafka.handler.sink.UpsertSinkHandler
  • 重要级别:中
  • 有效值 / 注意事项:使用可加载且具有可访问构造方法的 SinkHandler 类;非 DOCUMENTcouchbase.document.mode 会覆盖该值。

couchbase.document.mode

使用兼容模式选择文档、子文档或 N1QL 处理方式。
  • 类型string
  • 默认值DOCUMENT
  • 重要级别:中
  • 有效值 / 注意事项:可选 DOCUMENTSUBDOCUMENTN1QLN1QLSUBDOCUMENT 会覆盖 couchbase.sink.handler
  • 弃用状态:已弃用
  • 替代项couchbase.sink.handler

couchbase.retry.timeout

设置 Connector 重试失败写入的总时间窗口。
  • 类型string
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:使用非负整数加 mssmhd0 表示不重试。该配置不同于 Kafka Connect 的 errors.retry.timeout,失败时可能重放整个批次。

写入持久性

couchbase.durability

设置 Couchbase Server 6.5 或更高版本的增强持久性级别。
  • 类型string
  • 默认值NONE
  • 重要级别:中
  • 有效值 / 注意事项:可选 NONEMAJORITYMAJORITY_AND_PERSIST_TO_ACTIVEPERSIST_TO_MAJORITY。非 NONE 值不能与非 NONEcouchbase.persist.tocouchbase.replicate.to 同时使用。

couchbase.persist.to

设置旧版 Couchbase Server 的磁盘持久化要求。
  • 类型string
  • 默认值NONE
  • 重要级别:中
  • 有效值 / 注意事项:可选 NONEACTIVEONETWOTHREEFOUR,且要求必须能由集群拓扑满足。Couchbase Server 6.5 或更高版本优先使用 couchbase.durability

couchbase.replicate.to

设置旧版 Couchbase Server 的内存副本确认要求。
  • 类型string
  • 默认值NONE
  • 重要级别:中
  • 有效值 / 注意事项:可选 NONEONETWOTHREE,且要求必须能由集群拓扑满足。Couchbase Server 6.5 或更高版本优先使用 couchbase.durability

N1QL 处理器

couchbase.n1ql.operation

选择 N1qlSinkHandler 的更新方式。
  • 类型string
  • 默认值UPDATE
  • 重要级别:中
  • 有效值 / 注意事项:可选 UPDATEUPDATE_WHERE;仅适用于 N1qlSinkHandler

couchbase.n1ql.where.fields

指定构造 UPDATE_WHERE 条件使用的字段和可选字面量。
  • 类型list
  • 默认值:空列表
  • 重要级别:中
  • 有效值 / 注意事项:使用字段名或 field:literal 条目;选择 UPDATE_WHERE 时必须提供非空列表。

couchbase.n1ql.create.document

控制 N1QL UPDATE 模式是否通过 MERGE 创建不存在的文档。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项:仅适用于 N1qlSinkHandlerUPDATE 模式。

Analytics 处理器

couchbase.analytics.max.records.in.batch

设置每个 Analytics UPSERT 或 DELETE 语句批次的最大记录数。
  • 类型int
  • 默认值100
  • 重要级别:中
  • 有效值 / 注意事项:ConfigDef 接受任意 32 位整数;实际应使用正整数。该配置尚未承诺为稳定接口。

couchbase.analytics.max.size.in.batch

设置 Analytics UPSERT 批次中文档数据的最大总字节数。
  • 类型string
  • 默认值5m
  • 重要级别:中
  • 有效值 / 注意事项:使用非负整数加 bkmg;实际应使用正的大小限制。该配置尚未承诺为稳定接口。

couchbase.analytics.query.timeout

设置 Analytics 查询请求的客户端超时时间。
  • 类型string
  • 默认值5m
  • 重要级别:中
  • 有效值 / 注意事项:使用非负整数加 mssmhd,也可使用 0。该配置尚未承诺为稳定接口。

子文档处理器

couchbase.subdocument.path

指定 SubDocumentSinkHandler 修改的子文档路径。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:必须使用非空路径;普通值表示固定路径,以 / 开头的值表示从每条消息中提取动态路径的 JSON Pointer。

couchbase.subdocument.operation

选择子文档变更操作。
  • 类型string
  • 默认值UPSERT
  • 重要级别:中
  • 有效值 / 注意事项:可选 UPSERTARRAY_PREPENDARRAY_APPEND;仅适用于 SubDocumentSinkHandler

couchbase.subdocument.create.path

控制子文档变更是否创建缺失的父路径。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项:可选 truefalse;当 couchbase.subdocument.create.document=true 时,设为 false 不改变创建文档的行为。

couchbase.subdocument.create.document

控制子文档变更是否创建不存在的目标文档。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项true 使用 UPSERT 存储语义,false 使用 REPLACE 存储语义。

couchbase.create.document

为 N1QL 或子文档处理器设置创建文档行为的旧式兼容别名。
  • 类型boolean
  • 默认值:无固定默认值
  • 重要级别:中
  • 有效值 / 注意事项:仅在对应的当前配置未设置时复制到处理器配置。
  • 弃用状态:已弃用
  • 替代项couchbase.n1ql.create.documentcouchbase.subdocument.create.document

日志与指标

couchbase.log.redaction

设置 Couchbase 日志标记的脱敏级别。
  • 类型string
  • 默认值NONE
  • 重要级别:中
  • 有效值 / 注意事项:可选 NONEPARTIALFULL

couchbase.log.document.lifecycle

控制是否将文档生命周期消息提升到 INFO 级别。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:可选 truefalse;启用可能产生较高日志量。

couchbase.metrics.interval

设置 Connector 将指标快照写入日志的时间间隔。
  • 类型string
  • 默认值10m
  • 重要级别:中
  • 有效值 / 注意事项:使用非负整数加 mssmhd0 禁用指标日志。该配置尚未承诺为稳定接口。

SDK 扩展

couchbase.env.*

将 Couchbase Java SDK 属性传递给 ClusterEnvironment
  • 类型SDK-dependent
  • 默认值SDK default
  • 重要级别:中
  • 有效值 / 注意事项:在 couchbase.env. 后填写去掉 com. 前缀的 SDK 属性名;无效属性会导致 Task 启动失败。该模式不是 Connector ConfigDef 中的固定配置项。

最佳实践

使用稳定业务键更新同一文档

适用业务场景:首次接入持续变化的业务实体时,消息值中已有稳定且唯一的业务 ID,需要让重复投递和后续更新继续作用于同一 Couchbase 文档,而不是按 Kafka 记录坐标创建多个文档。 配置示例
关键说明${/orderId} 从消息值中的标量字段提取文档 ID。示例同时从存储正文移除该字段;如果业务仍需查询该字段,可以保留默认的 false。字段缺失、为 null 或不是标量时会回退到 Kafka 键或记录坐标,因此应在接入前保证业务 ID 完整且稳定。

将多个 Topic 路由到不同 Collection

适用业务场景:完成单一数据流接入后,需要让一个 Connector 消费多个业务 Topic,并将不同业务域的数据持续写入各自的 Couchbase Collection,同时复用连接与认证配置。 配置示例
关键说明:Java properties 会将 \u005B\u005D 分别解析为 [],因此示例实际使用 couchbase.default.collection[orders]couchbase.default.collection[payments]。按 Topic 的上下文配置覆盖未带后缀的默认 Collection;目标 Scope 与 Collection 必须提前创建。该路由只依据原始 Kafka Topic,不会根据分区或消息字段动态选择 Collection。

为生产写入设置有限重试和确认级别

适用业务场景:Connector 已能稳定写入,需要降低短暂网络或目标端故障造成的立即失败,同时要求 Couchbase 在达到指定副本确认条件后再把写入视为完成。 配置示例
关键说明2m 是示例重试窗口,应结合故障恢复目标和告警响应时间调整;Connector 不区分瞬时错误与永久数据错误,过长窗口可能延迟暴露坏记录。MAJORITY 需要目标集群支持且有足够副本可用。重试可能重新执行已部分成功的批次,这组配置不提供 Kafka Offset 与 Couchbase 写入之间的恰好一次保证;使用数组追加、条件更新或自定义非幂等处理器时,应额外评估重复副作用。

监控

监控内容

关注 Kafka Connect 通用健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试以及 Worker JVM 的堆内存、GC、线程和 CPU 信号;仅在部署启用了相应错误处理时关注 DLQ 活动。

导入 Grafana 大盘

确认 Kafka Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。

限制条件

  • Couchbase 写入与 Kafka Offset 提交不构成原子事务;写入成功但 Offset 提交前发生故障时,记录可能在恢复后重放,因此不能提供端到端恰好一次保证。
  • 不同 Task 之间没有共享写入顺序协调;来自不同 Kafka 分区的记录映射到同一文档 ID 时可以并发写入,不能保证跨分区、跨 Task 或跨文档的全局顺序。
  • 默认处理器执行整文档 upsert,而不是字段合并;相同文档 ID 的后续记录会用新的完整 JSON 覆盖目标文档。
  • 子文档数组追加和前置操作不是重放幂等操作;批次部分成功后重试可能再次插入相同元素。
  • Kafka Connect 的逐记录错误容忍和 DLQ 不会自动处理发生在 SinkTask.put 内的 Couchbase Handler 或写入错误。
  • couchbase.durability 为非 NONE 时,couchbase.persist.tocouchbase.replicate.to 必须均为 NONE

常见问题

为什么相同业务实体被写成了多个文档?

couchbase.document.id 未成功提取业务 ID,且 Kafka 记录键为空或类型不受支持时,Connector 会使用 Topic、分区和 Offset 生成回退 ID。检查消息是否始终包含稳定的标量业务字段或稳定 Kafka 键;需要从消息值提取时,配置例如 couchbase.document.id=${/orderId},并确认该字段没有缺失、null 或对象、数组值。

为什么记录被写入了错误的 Collection?

检查 couchbase.default.collection 是否使用 scope.collectionbucket.scope.collection 格式,并确认按 Topic 覆盖中的 Topic 名称与实际记录 Topic 完全一致。优先使用 couchbase.default.collection[<kafka-topic>],避免与已弃用的 couchbase.topic.to.collection 混用;该路由不会根据消息字段自动切换 Collection。

为什么启用 TLS 后 Task 无法连接 Couchbase?

检查 couchbase.enable.tls、证书主机名和 Worker 上的证书文件路径。PEM CA 使用 couchbase.trust.certificate.path,Java 信任库使用 couchbase.trust.store.path 及其密码;双向 TLS 还需配置客户端证书路径和密码。除非明确了解风险,不要关闭 couchbase.enable.hostname.verification

为什么配置了 Kafka Connect 错误容忍或 DLQ,Couchbase 写入错误仍使 Task 失败?

Kafka Connect 的逐记录错误处理主要覆盖 Converter 和 SMT 等阶段,Couchbase Handler 与目标写入发生在 SinkTask.put 内,不会自动被 errors.tolerance 跳过或写入 DLQ。检查 Task 日志中的具体写入错误,修复无效数据、权限、Keyspace 或连接问题;对于可恢复故障,可设置有限的 couchbase.retry.timeout,并配合告警避免永久错误长时间重试。

为什么增加 tasks.max 后吞吐没有提升或同一文档的更新顺序改变?

实际 Task 并行度受 Kafka Topic 分区数量和 Worker 分配限制,单个分区不会被拆给多个 Task。若不同分区中的记录使用同一 Couchbase 文档 ID,不同 Task 可以并发写入该文档,最终结果取决于目标端完成顺序。需要保持实体内顺序时,应让同一实体稳定落入同一 Kafka 分区,并结合目标端容量逐步调整 tasks.max

为什么子文档写入失败或意外创建了文档?

确认 couchbase.sink.handler 已选择 SubDocumentSinkHandler,并提供非空的 couchbase.subdocument.path。数组追加或前置要求目标路径满足 Couchbase 数组操作条件;只允许修改已有文档时,将 couchbase.subdocument.create.document 设为 falsecouchbase.subdocument.create.path 只控制缺失父路径,不能替代目标文档存在性控制。