跳到主要内容
版本:Next

Couchbase

Couchbase 接收器连接器

支持的引擎

Spark
Flink
SeaTunnel Zeta

主要特性

描述

将数据写入 Couchbase 集合。 每行数据以 JSON 文档形式存储。文档键由 primary-key 字段的值使用长度前缀规范编码构建 (格式为 <长度>:<值>,各分量以 # 分隔,例如 3:foo#3:bar)。 此编码不会产生碰撞:包含分隔符或特殊字符的值不会与其他不同的元组生成相同的键。 未配置时使用随机 UUID 作为文档键。

连接器支持:

  • Upsert 模式 — 插入或替换已有文档。
  • 批量刷写 — 在内存中缓冲数据,按行数或时间阈值刷写。
  • 重试机制 — 写入失败时采用线性退避重试(第 n 次重试等待 retry.interval × n 毫秒)。

支持的数据源信息

使用 Couchbase 连接器需要以下依赖。

数据源支持版本依赖
CouchbaseServer 7.x+下载

数据库依赖

运行作业前请安装连接器插件:

sh bin/install-plugin.sh ${version}

数据类型映射

SeaTunnel 数据类型Couchbase JSON 值
BOOLEANBoolean
TINYINT / SMALLINT / INTNumber (整数)
BIGINTNumber (长整数)
FLOAT / DOUBLENumber (浮点数)
DECIMALString (精确小数,如 "123.456")
STRINGString
DATE / TIME / TIMESTAMPString (ISO-8601)
BYTESString (Base64 编码)
ARRAYArray(元素递归转换)
MAPObject(键强制转为 String,值递归转换)
ROWObject(嵌套 JSON 文档)
NULLnull

接收器选项

名称类型是否必填默认值描述
connection.stringString-Couchbase 连接字符串,例如 couchbase://localhost
usernameString-Couchbase 用户名。
passwordString-Couchbase 密码。
bucketString-目标 Bucket 名称。
scopeString_defaultBucket 中的目标 Scope 名称。
collectionString-目标 Collection 名称。
ready.timeoutInteger30写入器初始化时等待目标 bucket 就绪的最长时间(秒),必须大于零。
primary-keyList<String>-用于构建文档键的字段名列表(长度前缀编码:<长度>:<值> 分量以 # 分隔)。未设置时使用随机 UUID。
upsert-enableBooleanfalse是否启用 Upsert(插入或替换)模式。为 false 时,重复键将报错。
buffer-flush.max-rowsInteger1000触发批量写入的最大缓冲行数。设为 -1 禁用。
retry.maxInteger3写入失败时的最大重试次数。
retry.intervalLong1000线性退避基础间隔(毫秒)。第 n 次重试等待 retry.interval × n 毫秒。

启动就绪等待

ready.timeout 控制写入器初始化时等待目标 bucket 就绪的时间,默认仍为 30 秒。 对于需要更长时间才能恢复可用的集群,可配置更大的正数,例如 ready.timeout = 60

该值的单位是秒,而不是毫秒。连接器不额外限制上限,应根据集群实际恢复时间选择足够的最小值。 当 bucket 持续不可用时,过大的值会延迟写入器初始化失败的报告。

Couchbase SDK 在此等待期间处理连接尝试,连接器不会额外添加启动重试循环。 retry.maxretry.interval 仍仅用于写入重试。此选项不会修改 SDK 单次操作的超时时间, 也不会修改引擎的作业启动超时时间。等待超时后,写入器初始化仍会失败并断开客户端连接。 增加等待时间不能修复无效凭据、错误地址或不存在的 bucket。

安全性

TLS / 加密传输

在生产环境中,请使用 couchbases:// 协议(注意末尾的 s)来启用 TLS 加密传输。 如需配置 CA 证书或自定义信任库,可通过 Couchbase Java SDK 的 ClusterEnvironment 进行设置:

sink {
Couchbase {
# 使用 couchbases://(末尾带 's')以启用 TLS 加密传输
connection.string = "couchbases://couchbase.example.com"
username = "seatunnel_writer"
password = "${env:COUCHBASE_PASSWORD}"
bucket = "my_bucket"
collection = "my_collection"
}
}

有关 TLS 配置、证书固定、客户端证书和密码套件选项,请参阅 Couchbase Java SDK — 安全连接

最小权限服务账户

生产环境中不应使用内置的 Administrator 账户。 请创建一个专用的 Couchbase 用户,并仅授予所需的最小权限:

  • 针对目标 Bucket/Scope/Collection 的 Data Writer 角色(仅插入场景)。
  • Data Reader + Data Writer 角色(Upsert 场景下可能需要读后写)。

凭据保护

请避免在 Job 配置文件中以明文存储密码。 SeaTunnel 支持加密配置值,详情请参阅 SeaTunnel 凭据加密文档, 了解如何在运行时动态替换密钥。

任务示例

简单示例 (仅供开发使用)

⚠️ 以下连接字符串和凭据仅适用于本地开发环境。 生产部署前请参阅上方的安全性章节。

sink {
Couchbase {
connection.string = "couchbase://127.0.0.1"
username = "Administrator"
password = "password"
bucket = "my_bucket"
collection = "my_collection"
}
}

使用 Upsert 和复合文档键 (仅供开发使用)

⚠️ 以下连接字符串和凭据仅适用于本地开发环境。 生产部署前请参阅上方的安全性章节。

sink {
Couchbase {
connection.string = "couchbase://127.0.0.1"
username = "Administrator"
password = "password"
bucket = "my_bucket"
scope = "_default"
collection = "my_collection"
primary-key = ["user_id", "order_id"]
upsert-enable = true
buffer-flush.max-rows = 500
retry.max = 5
retry.interval = 2000
}
}

定时刷新

该连接器支持在空闲时按计时器刷新缓冲区,即使尚未达到 buffer-flush.max-rows,也能定期发送已缓冲的记录。
此定时器由引擎驱动,而非连接器本身,目前仅 SeaTunnel Zeta 支持。

在 Job 的 env 块中设置 sink.flush.interval(毫秒)即可启用:

env {  
sink.flush.interval = 10000
}

在 Spark 和 Flink 上没有检查点之间的定时刷新:sink.flush.interval 是 Zeta 引擎的能力,
Spark/Flink 的 Sink 写入器上下文并未实现它。在这两个引擎上,缓存会在达到
buffer-flush.max-rows、检查点时(CouchbaseWriterprepareCommit() 中刷新)
以及写入器关闭时被刷新。如需降低 Spark 或 Flink 上检查点之间的延迟,
请相应调整 buffer-flush.max-rows

Change Log
ChangeCommitVersion