ScyllaDB Kafka Sink Connector 配置完全指南:属性详解、SSL/认证与一致性级别实战
2026/9/15 11:18:26 网站建设 项目流程

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/config

Standalone 模式(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.classString—(必填)Highio.connect.scylladb.ScyllaDbSinkConnector
scylladb.contact.pointsList[localhost]High多节点至少两个主机
scylladb.portInt9042Medium1–65535
scylladb.loadbalancing.localdcstring""High大小写敏感
scylladb.security.enabledBooleanfalseHigh需服务端认证配套
scylladb.usernameStringcassandraHigh需开启 security
scylladb.passwordPasswordcassandraHigh需开启 security
scylladb.compressionstringNONELowNONE/SNAPPY/LZ4(连接级)
scylladb.ssl.enabledbooleanfalseHigh需服务端 SSL 配套
scylladb.ssl.truststore.pathstring""MediumJava Truststore 路径
scylladb.ssl.truststore.passwordpassword[hidden]Medium
scylladb.ssl.providerstringJDKLowJDK/OPENSSL/OPENSSL_REFCNT
scylladb.keyspaceString—(必填)High相当于“数据库”
scylladb.keyspace.create.enabledBooleantrueHighfalse 且需重建时报错
scylladb.keyspace.replication.factorint3High≥1,单节点建议 1
scylladb.table.manage.enabledBooleantrueHighfalse 需表已存在
scylladb.table.create.compression.algorithmstringNONEMediumNONE/SNAPPY/LZ4/DEFLATE(表级)
scylladb.offset.storage.tableStringkafka_connect_offsetsLow偏移量存储表
scylladb.consistency.levelStringLOCAL_QUORUMHigh见上文 11 种取值
scylladb.deletes.enabledbooleantrueHighnull 值记录触发删除
scylladb.execute.timeout.msLong30000Low语句执行超时
scylladb.ttlIntnullMedium不配则无 TTL
scylladb.offset.storage.table.enableBooleantrueMediumfalse 可能重复写入
scylladb.max.batch.size.kbint5High需与 scylla.yaml 阈值联动
tasks.maxint—(必填)High并行度
topicslist—(必填)High消费的 Topic 列表
confluent.topic.bootstrap.serverslist—(必填)Highhost:port列表

十、排错与日志定位

  • Connector 日志(Confluent 平台)

    confluent local log <service> -- [<argument>] --path <path-to-confluent>
  • ScyllaDB 日志(Docker 部署)

    docker logs some-scylla | tail
  • 常见问题速查:

    1. NoHostAvailableException/ 连接失败 → 检查scylladb.contact.pointsscylladb.port(Docker 场景用docker ps确认映射端口);
    2. Unauthorized→ 检查scylladb.security.enabled、用户名/密码,以及服务端authenticator配置;
    3. Keyspace 创建报错 → 检查scylladb.keyspace.create.enabled与 RF 设置(单节点把 RF 调为 1);
    4. 批次过大被服务端拒绝 → 按第六节的比例关系同步调整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_kbbatch_size_fail_threshold_in_kbauthenticator
  • 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询