ScyllaDB Kafka Sink Connector 配置完全指南:属性详解、SSL/认证与一致性级别实战
【免费下载链接】scylladbNoSQL data store using the Seastar framework, compatible with Apache Cassandra and Amazon DynamoDB项目地址: https://gitcode.com/GitHub_Trending/sc/scylladb
导读
本文以 ScyllaDB 官方文档《Kafka Sink Connector Configuration》(docs/using-scylla/integrations/sink-config.rst)为骨架,系统梳理 ScyllaDB Kafka Sink Connector 的每一个配置属性:连接参数、SSL 与安全认证、Keyspace/Table 管理、写入行为、偏移量存储以及 Confluent 平台相关配置。读完本文,你将能够独立编写一份可直接运行的 Connector 配置文件(分布式 JSON 或独立 properties 格式),并理解每个参数背后的实现逻辑与调优要点——例如scylladb.max.batch.size.kb为何必须与 ScyllaDB 服务端的批次阈值参数联动调整。
定位说明:本文聚焦于 KafkaSinkConnector(把 Kafka 记录写入 ScyllaDB)。快速上手流程(Docker 启动 ScyllaDB、Confluent 安装、
mvn clean install构建)见配套文档 Kafka Sink Connector Quickstart,本文不重复展开。
一、Connector 类名:一切配置的起点
Sink Connector 的入口是一个 Java 类,使用前必须通过connector.class属性显式声明:
connector.class=io.connect.scylladb.ScyllaDbSinkConnector在 Confluent 平台中,可以通过 REST API 校验该插件是否已被 Connect 的插件加载器正确识别:
curl -sS localhost:8083/connector-plugins | jq .[].class | grep ScyllaDbSinkConnector正常输出为io.connect.scylladb.ScyllaDbSinkConnector,此时才说明插件安装成功,可以继续加载 Connector 实例(安装步骤参见 kafka-connector.rst 中的 “Add Sink Connector plugin” 小节)。
二、Connection:连接参数
scylladb.contact.points
指定要连接的 ScyllaDB 节点主机列表。ScyllaDB 节点使用这份主机列表互相发现并学习整个 Ring 的拓扑结构。
- Type: List
- Importance: High
- Default Value:
[localhost] - 单节点场景下默认值即可;多节点集群必须修改,且出于高可用考虑,规模较大的集群至少应配置两个主机;
- 若使用 Docker 镜像运行 ScyllaDB,需填写该容器对外映射的宿主机地址。
scylladb.port
ScyllaDB 节点监听的端口。
- Type: Int
- Importance: Medium
- Default Value:
9042 - Valid Values:
ValidPort{start=1, end=65535} - 使用 Docker 镜像时,用
docker ps查看容器实际映射出来的端口并填入。
scylladb.loadbalancing.localdc
指定与 Connector 运行机器同处的本地数据中心名称(大小写敏感)。设置后,驱动会把流量优先发往本地 DC,避免跨数据中心访问延迟。
- Type: string
- Default:
"" - Importance: High
scylladb.security.enabled
启用安全模式后再加载 Sink Connector 并连接 ScyllaDB。
- Type: Boolean
- Importance: High
- Default Value:
False
需要说明的是:ScyllaDB 服务端默认使用
AllowAllAuthenticator(conf/scylla.yaml 第 250 行附近,默认不开启认证)。Connector 侧的scylladb.security.enabled是客户端行为开关,两者需要配套:服务端启用PasswordAuthenticator时,客户端必须同时打开此开关并提供用户名/密码。
scylladb.username
连接 ScyllaDB 使用的用户名。使用该参数时必须同时设置scylladb.security.enabled = true。
- Type: String
- Importance: High
- Default Value:
cassandra
scylladb.password
连接 ScyllaDB 使用的密码,同样要求scylladb.security.enabled = true。
- Type: Password
- Importance: High
- Default Value:
cassandra
scylladb.compression
连接 ScyllaDB 时使用的压缩算法,用于压缩客户端与节点之间的传输流量。
- Type: string
- Default:
NONE - Valid Values:
[NONE, SNAPPY, LZ4] - Importance: Low
- 注意:这是客户端连接协议级的压缩,与后文表级压缩
scylladb.table.create.compression.algorithm(SSTable 压缩)是两回事,不要混淆。
scylladb.ssl.enabled
连接 ScyllaDB 时是否启用 SSL。
- Type: boolean
- Default:
false - Importance: High
- 开启后需配合下面的 SSL 子配置项(Truststore)使用;同时要求 ScyllaDB 服务端已开启
native_transport_port_ssl对应的 SSL 监听。
三、SSL:Java Truststore 配置
启用scylladb.ssl.enabled=true后,由以下三个属性定义客户端信任链:
scylladb.ssl.truststore.path
Java Truststore 文件的路径。
- Type: string
- Default:
"" - Importance: medium
scylladb.ssl.truststore.password
访问 Java Truststore 的密码。
- Type: password
- Default:
[hidden] - Importance: medium
scylladb.ssl.provider
连接 ScyllaDB 时使用的 SSL Provider。
- Type: string
- Default:
JDK - Valid Values:
[JDK, OPENSSL, OPENSSL_REFCNT] - Importance: low
JDK使用 JSSE 默认实现;OPENSSL/OPENSSL_REFCNT走 Netty 的 OpenSSL 原生实现(后者额外维护引用计数),通常吞吐更高,但需要运行环境具备对应的 native 依赖。
四、Keyspace:目标库管理
scylladb.keyspace
要写入的 Keyspace 名称。Keyspace 相当于 ScyllaDB 集群中的“数据库”,表与数据都归属其下。
- Type: String
- Importance: High
- 无默认值,必须显式指定(Quickstart 示例中为
test)。
scylladb.keyspace.create.enabled
决定当 Keyspace 不存在时,Connector 是否自动创建它。
- Type: Boolean
- Importance: High
- Default Value:
true - ⚠️注意:如果 Keyspace 已存在且 Connector 需要重建(例如发生冲突),而该参数为
false,将会直接报错。因此生产中若已手工建好 Keyspace,应结合实际情况决定是否关闭自动创建。
scylladb.keyspace.replication.factor
Connector 自动创建 Keyspace 时使用的复制因子(Replication Factor,RF)。RF 等于数据(行与分区)被复制的节点份数,RF=N 表示数据复制到 N 个节点。
- Type: int
- Default:
3 - Valid Values:
[1,...] - Importance: High
- 单节点开发环境记得把 RF 调整为 1,否则写入可能因无法满足 RF 而失败。
五、Table:表管理与偏移量存储
scylladb.table.manage.enabled
决定 Connector 是否负责管理(创建/维护)目标表。
- Type: Boolean
- Importance: High
- Default Value:
true - 若为
false,则要求表已预先存在且结构匹配,Connector 只做写入。
scylladb.table.create.compression.algorithm
创建表时使用的表级(SSTable)压缩算法。
- Type: string
- Default:
NONE - Valid Values:
[NONE, SNAPPY, LZ4, DEFLATE] - Importance: medium
- 注意此处的取值范围比连接级压缩(
scylladb.compression)多一个DEFLATE。生产环境通常建议至少启用 LZ4 以降低磁盘占用。
scylladb.offset.storage.table
在 ScyllaDB Keyspace 内用于存储“已从 Kafka 读取到的偏移量(offset)”的表名。该表只需启用一次即可实现向 ScyllaDB 的可靠投递。
- Type: String
- Importance: Low
- Default:
kafka_connect_offsets
六、Write:写入行为控制
scylladb.consistency.level
写入 ScyllaDB 时要求的一致性级别(Consistency Level,CL)。CL 决定一个读写操作被判定成功前,集群中必须有多少个副本进行确认。
- Type: String
- Importance: High
- Default Value:
LOCAL_QUORUM - Valid Values:
ANY, ONE, TWO, THREE, QUORUM, ALL, LOCAL_QUORUM, EACH_QUORUM, SERIAL, LOCAL_SERIAL, LOCAL_ONE
从源码角度印证:ScyllaDB 服务端将一致性级别建模为密集枚举db::consistency_level(db/consistency_level_type.hh 第 20–32 行),取值依次为ANY, ONE, TWO, THREE, QUORUM, ALL, LOCAL_QUORUM, EACH_QUORUM, SERIAL, LOCAL_SERIAL, LOCAL_ONE,与本文档的合法值列表完全一致。其实现语义(db/consistency_level.cc)可进一步佐证:
LOCAL_QUORUM(默认值)要求本地 DC内满足法定数量副本确认,若本地 DC 的 RF 为 0 或存活副本不足,会抛出unavailable_exception(第 215–225 行);ANY仅要求本地节点接受写入(可先落 hint,第 204–206 行),代价是最低一致性;EACH_QUORUM要求每个 DC都满足 QUORUM,多 DC 场景下代价较高(第 227 行起)。
scylladb.deletes.enabled
决定 Connector 是否处理删除操作。Kafka 记录值为null时,将按该记录 Key 中的主键删除 ScyllaDB 中对应的记录。
- Type: boolean
- Default:
true - Importance: High
scylladb.execute.timeout.ms
执行一条 ScyllaDB 语句的超时时间。
- Type: Long
- Importance: Low
- Default Value:
30000(毫秒)
scylladb.ttl
数据在 ScyllaDB 中的保留期限(Time To Live)。超过该间隔后,ScyllaDB 会自动清除这些记录。若不配置,Sink Connector 执行插入时将不带 TTL 设置(数据永久保留,除非显式删除或服务端另有策略)。
- Type: Int
- Importance: Medium
- Default Value:
null(不生效)
scylladb.offset.storage.table.enable
true:Kafka 消费偏移量存入 ScyllaDB 表(表名见scylladb.offset.storage.table);false:Connector 跳过向 ScyllaDB 写偏移量,任务重启时可能产生重复写入(at-least-once 语义)。Type: Boolean
Importance: Medium
Default Value:
True
scylladb.max.batch.size.kb
单个批次(由若干 ScyllaDB 操作组成)的最大大小,单位 KB。该值与 ScyllaDB 服务端配置强相关,文档明确要求:
应等于
scylla.yaml中配置的batch_size_warn_threshold_in_kb,并为batch_size_fail_threshold_in_kb的 1/10。默认值 5KB,任何修改都必须同步修改scylla.yaml。
- Type: int
- Default:
5 - Valid Values:
[1,...] - Importance: High
在仓库自带的 conf/scylla.yaml 第 229–234 行可以看到服务端默认值:
# Log WARN on any batch size exceeding this value. 128 kiB per batch by default. # Caution should be taken on increasing the size of this threshold as it can lead to node instability. batch_size_warn_threshold_in_kb: 128 # Fail any multiple-partition batch exceeding this value. 1 MiB (8x warn threshold) by default. batch_size_fail_threshold_in_kb: 1024即服务端默认 128KB 告警、1024KB(1MiB)拒绝。Connector 默认 5KB 远低于告警线,属于保守取值;若你希望放大批次提升吞吐,应同步上调服务端的两个阈值(并注意文档“增大阈值可能导致节点不稳定”的告诫),再按“告警值 = 1/10 × 失败值”的比例关系设置scylladb.max.batch.size.kb。
七、Confluent Platform 配置
这部分是 Confluent Kafka Connect 平台层面的通用属性,与具体 Connector 实现无关但必填:
tasks.max
Connector 使用的最大任务(Task)数,用于实现并行度。
- Type: int
- Importance: High
topics
要从哪些 Kafka Topic 消费数据并写入 ScyllaDB(列表)。
- Type: list
- Importance: High
- 示例:
topics=topic1,topic2,topic3;最终写入的目标表名通常取自 Topic 名(Quickstart 中select * from demo.example即对应 topicexample)。
confluent.topic.bootstrap.servers
用于建立到 Kafka 集群初始连接的 host/port 列表(主要用于许可/执照验证)。格式为host1:port1,host2:port2,…。集群中的其余服务器会从初始连接自动发现;因此该列表不必包含全部服务器,但建议至少多写一个以防某台宕机。
- Type: list
- Importance: High
八、完整配置示例:可直接落地的最小可用配置
Distributed 模式(JSON,通过 REST API 加载)
{ "name" : "scylladb-sink-connector", "config" : { "connector.class" : "io.connect.scylladb.ScyllaDbSinkConnector", "tasks.max" : "1", "topics" : "topic1,topic2,topic3", "scylladb.contact.points" : "scylladb-hosts", "scylladb.keyspace" : "test", "key.converter" : "org.apache.kafka.connect.json.JsonConverter", "value.converter" : "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable" : "true", "value.converter.schemas.enable" : "true", "transforms" : "createKey", "transforms.createKey.fields" : "[field-you-want-as-primary-key-in-scylla]", "transforms.createKey.type" : "org.apache.kafka.connect.transforms.ValueToKey" } }加载与更新:
# 创建 Connector curl -s -X POST -H 'Content-Type: application/json' --data @kafka-connect-scylladb.json http://localhost:8083/connectors # 更新已有 Connector 配置 curl -s -X PUT -H 'Content-Type: application/json' --data @kafka-connect-scylladb.json http://localhost:8083/connectors/scylladb/configStandalone 模式(properties 文件)
connector.class=io.connect.scylladb.ScyllaDbSinkConnector tasks.max=1 topics=topic1,topic2,topic3 scylladb.contact.points=cassandra scylladb.keyspace=test key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=true value.converter.schemas.enable=true transforms=createKey transforms.createKey.fields=[field-you-want-as-primary-key-in-scylla] transforms.createKey.type=org.apache.kafka.connect.transforms.ValueToKey加载命令:
confluent local load scylladb-sink-conector -- -d scylladb-sink-connector.properties提示:
transforms.createKey通过ValueToKey转换把 payload 中的某个字段提升为记录 Key(即 ScyllaDB 主键),这是“无 Key 记录写入 ScyllaDB”场景下的关键配置。若使用 Avro 序列化,把两个 converter 换成io.confluent.connect.avro.AvroConverter并配置schema.registry.url=http://localhost:8081即可。
启用认证与 SSL 的配置
Distributed 模式:
{ "name" : "scylladbSinkConnector", "config" : { "connector.class" : "io.connect.scylladb.ScyllaDbSinkConnector", "tasks.max" : "1", "topics" : "topic1,topic2,topic3", "scylladb.contact.points" : "cassandra", "scylladb.keyspace" : "test", "scylladb.security.enabled" : "true", "scylladb.username" : "example", "scylladb.password" : "password" } }Standalone 模式:
connector.class=io.connect.scylladb.ScyllaDbSinkConnector tasks.max=1 topics=topic1,topic2,topic3 scylladb.contact.points=cassandra scylladb.keyspace=test scylladb.ssl.enabled=true scylladb.username=example scylladb.password=password前提约束:客户端开启
scylladb.security.enabled前,ScyllaDB 服务端须已把 conf/scylla.yaml 中的authenticator从默认的AllowAllAuthenticator切换为PasswordAuthenticator(并相应提升system_authkeyspace 的 RF);开启 SSL 前,服务端须已配置好证书与 SSL 监听端口。
验证写入结果
写入数据后(Quickstart 中通过kafka-avro-console-producer/kafka-console-producer生产记录),在 ScyllaDB 侧用 cqlsh 验证:
cqlsh> SELECT * FROM test.topic1;九、参数速查表
| 配置项 | 类型 | 默认值 | 重要性 | 合法值/备注 |
|---|---|---|---|---|
connector.class | String | —(必填) | High | io.connect.scylladb.ScyllaDbSinkConnector |
scylladb.contact.points | List | [localhost] | High | 多节点至少两个主机 |
scylladb.port | Int | 9042 | Medium | 1–65535 |
scylladb.loadbalancing.localdc | string | "" | High | 大小写敏感 |
scylladb.security.enabled | Boolean | false | High | 需服务端认证配套 |
scylladb.username | String | cassandra | High | 需开启 security |
scylladb.password | Password | cassandra | High | 需开启 security |
scylladb.compression | string | NONE | Low | NONE/SNAPPY/LZ4(连接级) |
scylladb.ssl.enabled | boolean | false | High | 需服务端 SSL 配套 |
scylladb.ssl.truststore.path | string | "" | Medium | Java Truststore 路径 |
scylladb.ssl.truststore.password | password | [hidden] | Medium | — |
scylladb.ssl.provider | string | JDK | Low | JDK/OPENSSL/OPENSSL_REFCNT |
scylladb.keyspace | String | —(必填) | High | 相当于“数据库” |
scylladb.keyspace.create.enabled | Boolean | true | High | false 且需重建时报错 |
scylladb.keyspace.replication.factor | int | 3 | High | ≥1,单节点建议 1 |
scylladb.table.manage.enabled | Boolean | true | High | false 需表已存在 |
scylladb.table.create.compression.algorithm | string | NONE | Medium | NONE/SNAPPY/LZ4/DEFLATE(表级) |
scylladb.offset.storage.table | String | kafka_connect_offsets | Low | 偏移量存储表 |
scylladb.consistency.level | String | LOCAL_QUORUM | High | 见上文 11 种取值 |
scylladb.deletes.enabled | boolean | true | High | null 值记录触发删除 |
scylladb.execute.timeout.ms | Long | 30000 | Low | 语句执行超时 |
scylladb.ttl | Int | null | Medium | 不配则无 TTL |
scylladb.offset.storage.table.enable | Boolean | true | Medium | false 可能重复写入 |
scylladb.max.batch.size.kb | int | 5 | High | 需与 scylla.yaml 阈值联动 |
tasks.max | int | —(必填) | High | 并行度 |
topics | list | —(必填) | High | 消费的 Topic 列表 |
confluent.topic.bootstrap.servers | list | —(必填) | High | host:port列表 |
十、排错与日志定位
Connector 日志(Confluent 平台):
confluent local log <service> -- [<argument>] --path <path-to-confluent>ScyllaDB 日志(Docker 部署):
docker logs some-scylla | tail常见问题速查:
NoHostAvailableException/ 连接失败 → 检查scylladb.contact.points与scylladb.port(Docker 场景用docker ps确认映射端口);Unauthorized→ 检查scylladb.security.enabled、用户名/密码,以及服务端authenticator配置;- Keyspace 创建报错 → 检查
scylladb.keyspace.create.enabled与 RF 设置(单节点把 RF 调为 1); - 批次过大被服务端拒绝 → 按第六节的比例关系同步调整
scylladb.max.batch.size.kb与 conf/scylla.yaml 中的batch_size_warn_threshold_in_kb/batch_size_fail_threshold_in_kb。
延伸阅读
- Kafka Sink Connector Quickstart:Docker 启动 ScyllaDB、手动构建 Connector、加载与验证的完整上手流程
- Shard-Aware Kafka Connector for ScyllaDB:Kafka 集成系列文档入口
- 一致性级别服务端实现:db/consistency_level_type.hh、db/consistency_level.cc
- 服务端批次阈值与认证配置:conf/scylla.yaml(
batch_size_warn_threshold_in_kb、batch_size_fail_threshold_in_kb、authenticator) - ScyllaDB 原生 CDC Source Connector 文档:scylla-cdc-source-connector.rst(如需把 ScyllaDB 变更流回 Kafka,可参考同一集成体系)
【免费下载链接】scylladbNoSQL data store using the Seastar framework, compatible with Apache Cassandra and Amazon DynamoDB项目地址: https://gitcode.com/GitHub_Trending/sc/scylladb
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考