Couchbase
Couchbase 接收器连接器
支持的引擎
Spark
Flink
SeaTunnel Zeta
主要特性
描述
将数据写入 Couchbase 集合。
每行数据以 JSON 文档形式存储。文档键由 primary-key 字段的值使用长度前缀规范编码构建
(格式为 <长度>:<值>,各分量以 # 分隔,例如 3:foo#3:bar)。
此编码不会产生碰撞:包含分隔符或特殊字符的值不会与其他不同的元组生成相同的键。
未配置时使用随机 UUID 作为文档键。
连接器支持:
- Upsert 模式 — 插入或替换已有文档。
- 批量刷写 — 在内存中缓冲数据,按行数或时间阈值刷写。
- 重试机制 — 写入失败时采用线性退避重试(第 n 次重试等待
retry.interval × n毫秒)。
支持的数据源信息
使用 Couchbase 连接器需要以下依赖。
| 数据源 | 支持版本 | 依赖 |
|---|---|---|
| Couchbase | Server 7.x+ | 下载 |
数据库依赖
运行作业前请安装连接器插件:
sh bin/install-plugin.sh ${version}
数据类型映射
| SeaTunnel 数据类型 | Couchbase JSON 值 |
|---|---|
| BOOLEAN | Boolean |
| TINYINT / SMALLINT / INT | Number (整数) |
| BIGINT | Number (长整数) |
| FLOAT / DOUBLE | Number (浮点数) |
| DECIMAL | String (精确小数,如 "123.456") |
| STRING | String |
| DATE / TIME / TIMESTAMP | String (ISO-8601) |
| BYTES | String (Base64 编码) |
| ARRAY | Array(元素递归转换) |
| MAP | Object(键强制转为 String,值递归转换) |
| ROW | Object(嵌套 JSON 文档) |
| NULL | null |
接收器选项
| 名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| connection.string | String | 是 | - | Couchbase 连接字符串,例如 couchbase://localhost。 |
| username | String | 是 | - | Couchbase 用户名。 |
| password | String | 是 | - | Couchbase 密码。 |
| bucket | String | 是 | - | 目标 Bucket 名称。 |
| scope | String | 否 | _default | Bucket 中的目标 Scope 名称。 |
| collection | String | 是 | - | 目标 Collection 名称。 |
| ready.timeout | Integer | 否 | 30 | 写入器初始化时等待目标 bucket 就绪的最长时间(秒),必须大于零。 |
| primary-key | List<String> | 否 | - | 用于构建文档键的字段名列表(长度前缀编码:<长度>:<值> 分量以 # 分隔)。未设置时使用随机 UUID。 |
| upsert-enable | Boolean | 否 | false | 是否启用 Upsert(插入或替换)模式。为 false 时,重复键将报错。 |
| buffer-flush.max-rows | Integer | 否 | 1000 | 触发批量写入的最大缓冲行数。设为 -1 禁用。 |
| retry.max | Integer | 否 | 3 | 写入失败时的最大重试次数。 |
| retry.interval | Long | 否 | 1000 | 线性退避基础间隔(毫秒)。第 n 次重试等待 retry.interval × n 毫秒。 |
启动就绪等待
ready.timeout 控制写入器初始化时等待目标 bucket 就绪的时间,默认仍为 30 秒。
对于需要更长时间才能恢复可用的集群,可配置更大的正数,例如 ready.timeout = 60。
该值的单位是秒,而不是毫秒。连接器不额外限制上限,应根据集群实际恢复时间选择足够的最小值。 当 bucket 持续不可用时,过大的值会延迟写入器初始化失败的报告。
Couchbase SDK 在此等待期间处理连接尝试,连接器不会额外添加启动重试循环。
retry.max 和 retry.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、检查点时(CouchbaseWriter在prepareCommit()中刷新)
以及写入器关闭时被刷新。如需降低 Spark 或 Flink 上检查点之间的延迟,
请相应调整buffer-flush.max-rows。
Change Log
| Change | Commit | Version |
|---|