Skip to main content

概述

ScyllaDB Sink Connector 将 Kafka Topic 中的记录写入 ScyllaDB。Connector 使用 Kafka 记录的 key 和 value 组成目标表的主键与列,并根据 Topic 名称派生表名:Topic 中的句点和连字符会替换为下划线。默认情况下,Connector 可按记录的 Connect Schema 创建或补充表结构,也可以针对每个 Topic 配置列映射、一致性级别、TTL 和删除处理。 记录的 key 和 value 应为对象形态(Connect Struct 或 Map)。非空 value 用于写入;启用删除处理时,null value 可按记录 key 对目标表执行删除。Kafka Connect 负责订阅 Topic、任务分配和 Offset 提交,Connector 负责 ScyllaDB 会话、表结构处理和写入请求。

前置条件

  • 准备可由 Connect Worker 访问的 ScyllaDB 集群、目标 keyspace,以及在关闭自动建库时预先创建的 keyspace。
  • 使用结构化的 Kafka key 和 value;无 Schema 的 Map 数据需要预先创建兼容的 ScyllaDB 表,因为 Connector 不会从无 Schema 数据推断或修改表结构。
  • 启用 TLS 时准备可读的 truststore 或 keystore 文件及其密码;启用认证时同时提供用户名和密码。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka 和 ScyllaDB,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector
<topic-name><scylladb-host><keyspace-name> 替换为实际 Kafka Topic、ScyllaDB 地址和 keyspace。不要同时设置 topics.regex。记录 key 和 value 应由 Worker 的 Converter 转换为 Struct 或 Map。默认会创建 keyspace 和表,使用本地 QUORUM、启用删除处理和 ScyllaDB Offset 表。

配置

ScyllaDB 连接

scylladb.contact.points

ScyllaDB 联系点,可使用逗号分隔的地址列表,也可使用包含地址映射的 JSON。
  • 类型string
  • 默认值localhost
  • 重要级别:高
  • 有效值 / 注意事项:地址必须可由 Connect Worker 访问;Connector 校验时会尝试建立 ScyllaDB 会话。

scylladb.port

公共联系点使用的 ScyllaDB 端口。
  • 类型int
  • 默认值9042
  • 重要级别:中
  • 有效值 / 注意事项165535;JSON 私有地址映射可为每个地址指定自己的端口。

scylladb.loadbalancing.localdc

ScyllaDB 驱动使用的本地数据中心名称。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:区分大小写;为空时不显式指定本地数据中心。

scylladb.security.enabled

是否启用 ScyllaDB 用户名和密码认证。
  • 类型boolean
  • 默认值false
  • 重要级别:高
  • 有效值 / 注意事项:设为 true 时必须同时设置 scylladb.usernamescylladb.password

scylladb.username

ScyllaDB 认证用户名。
  • 类型string
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:必须与 scylladb.password 成对设置;启用认证时必填。

scylladb.password

ScyllaDB 认证密码。
  • 类型password
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:必须与 scylladb.username 成对设置;不要写入日志或文档示例。

scylladb.compression

ScyllaDB 驱动协议压缩方式。
  • 类型string
  • 默认值none
  • 重要级别:低
  • 有效值 / 注意事项nonelz4snappy,值使用小写。

TLS

scylladb.ssl.enabled

是否启用 ScyllaDB TLS 连接。
  • 类型boolean
  • 默认值false
  • 重要级别:高
  • 有效值 / 注意事项:设为 true 后按需设置 truststore、keystore、密码、密码套件和主机名校验配置。

scylladb.ssl.truststore.path

TLS truststore 文件路径。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:仅在 scylladb.ssl.enabled=true 时使用;设置后文件必须可读。

scylladb.ssl.truststore.password

TLS truststore 密码。
  • 类型password
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:仅在 scylladb.ssl.enabled=true 时使用;属于敏感值。

scylladb.ssl.keystore.path

TLS keystore 文件路径。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:仅在 scylladb.ssl.enabled=true 时使用;设置后文件必须可读。

scylladb.ssl.keystore.password

TLS keystore 密码。
  • 类型password
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:仅在 scylladb.ssl.enabled=true 时使用;属于敏感值。

scylladb.ssl.cipherSuites

TLS 允许使用的密码套件列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:高
  • 有效值 / 注意事项:仅在 scylladb.ssl.enabled=true 时使用;空列表表示不改变驱动的密码套件选择。

scylladb.ssl.hostname.verification

是否启用 TLS 主机名校验。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅在 scylladb.ssl.enabled=true 时使用。

写入行为

scylladb.consistency.level

写入 ScyllaDB 时使用的一致性级别。
  • 类型string
  • 默认值LOCAL_QUORUM
  • 重要级别:高
  • 有效值 / 注意事项ANYONETWOTHREEQUORUMALLLOCAL_QUORUMEACH_QUORUMSERIALLOCAL_SERIALLOCAL_ONE;匹配的 Topic 级配置会覆盖此值。

scylladb.deletes.enabled

是否将 null value 作为删除候选,并按记录 key 删除目标行。
  • 类型boolean
  • 默认值true
  • 重要级别:高
  • 有效值 / 注意事项:删除需要记录 key 包含目标表的全部主键字段;匹配的 Topic 级配置会覆盖此值。

scylladb.execute.timeout.ms

等待异步 ScyllaDB 操作完成的超时时间。
  • 类型long
  • 默认值30000
  • 重要级别:低
  • 有效值 / 注意事项:单位为毫秒,必须大于等于 00 也会被接受。

scylladb.ttl

插入语句使用的默认 TTL。
  • 类型int
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:为 null 时不添加 USING TTL;匹配的 Topic 级 ttlSeconds 会覆盖此值。

behavior.on.error

记录校验或语句构造发生 DataExceptionNullPointerException 时的处理方式。
  • 类型string
  • 默认值FAIL
  • 重要级别:中
  • 有效值 / 注意事项FAIL 抛出 Connect 异常,LOG 记录并继续,IGNORE 以跟踪级别记录并继续;不替代 Kafka Connect 的 errors.tolerance

Keyspace 与表

scylladb.keyspace

Connector 使用的目标 keyspace。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必填;启用自动建库时用于创建 keyspace,关闭自动建库时必须已存在。
  • 必填:是

scylladb.keyspace.create.enabled

是否自动创建不存在的 keyspace。
  • 类型boolean
  • 默认值true
  • 重要级别:高
  • 有效值 / 注意事项:设为 false 时,配置的 keyspace 必须已存在。

scylladb.keyspace.replication.factor

自动创建 keyspace 时使用的复制因子。
  • 类型int
  • 默认值3
  • 重要级别:高
  • 有效值 / 注意事项:必须大于等于 1;仅在 scylladb.keyspace.create.enabled=true 时生效。

scylladb.table.manage.enabled

是否由 Connector 创建或调整目标表结构。
  • 类型boolean
  • 默认值true
  • 重要级别:高
  • 有效值 / 注意事项:关闭后,目标表和所需列必须预先准备;不从无 Schema 数据生成 DDL。

scylladb.table.create.compression.algorithm

创建或调整表时使用的压缩算法。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项SnappyCompressorLZ4CompressorDeflateCompressornone;仅在表管理启用时影响创建或调整。

scylladb.offset.storage.table

保存 Connector ScyllaDB Offset 的表名。
  • 类型string
  • 默认值kafka_connect_offsets
  • 重要级别:低
  • 有效值 / 注意事项:仅在 scylladb.offset.storage.table.enable=true 时使用。

scylladb.offset.storage.table.enable

是否在 ScyllaDB 中创建、读取和写入 Offset 表。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项:关闭后跳过 ScyllaDB Offset 表;Kafka Connect 仍按其 Worker 机制管理 Kafka Offset。

Topic 到表映射

topic.<topic>.<keyspace>.<table>.mapping

为指定 Topic 派生的表配置列映射。
  • 类型string
  • 默认值null
  • 重要级别:未声明
  • 有效值 / 注意事项:使用逗号分隔的 列名=key.<字段>value.<字段>header.<字段>;也可映射特殊目标 __ttl__timestamp。Topic、keyspace 和 table 片段必须符合动态配置名称格式。

topic.<topic>.<keyspace>.<table>.consistencyLevel

为指定 Topic 派生的表覆盖写入一致性级别。
  • 类型string
  • 默认值:继承 scylladb.consistency.level
  • 重要级别:未声明
  • 有效值 / 注意事项:使用 ScyllaDB 支持的一致性级别;仅覆盖匹配 Topic 配置。

topic.<topic>.<keyspace>.<table>.ttlSeconds

为指定 Topic 派生的表覆盖默认 TTL。
  • 类型int
  • 默认值:继承 scylladb.ttl
  • 重要级别:未声明
  • 有效值 / 注意事项:按整数解析;空值继承 Connector 级 TTL。使用时应根据业务保留期设置。

topic.<topic>.<keyspace>.<table>.deletesEnabled

为指定 Topic 派生的表覆盖删除处理开关。
  • 类型boolean
  • 默认值:继承 scylladb.deletes.enabled
  • 重要级别:未声明
  • 有效值 / 注意事项:仅接受 truefalse,不区分大小写。

Kafka Connect Sink 框架

connector.class

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

tasks.max

允许 Connector 创建的最大 Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:必须大于等于 1;实际并行度还取决于 Kafka 分区分配和 Worker 调度。

topics

要消费的 Kafka Topic 列表。
  • 类型list
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:与 topics.regex 互斥,必须二选一;不能包含 DLQ Topic。

topics.regex

按 Java Pattern 语法匹配要消费的 Kafka Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:与 topics 互斥,必须二选一;不能匹配 DLQ Topic。

transforms

按顺序执行的 Kafka Connect SMT 列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:低
  • 有效值 / 注意事项:每个别名都需要对应的 transforms.<alias>.type;转换后的 Topic 名称用于 Topic 到表的映射查找。

predicates

供 SMT 条件判断使用的 Predicate 别名列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:低
  • 有效值 / 注意事项:仅在配置的 SMT 引用 Predicate 时使用;别名必须唯一。

errors.tolerance

Kafka Connect 框架处理错误的容忍范围。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项noneall;与 behavior.on.error 分开生效,主要影响转换、SMT 和错误报告阶段。

errors.retry.timeout

Kafka Connect 框架重试失败操作的总时长。
  • 类型long
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:单位为毫秒;-1 表示持续重试。它不改变 scylladb.execute.timeout.ms

errors.deadletterqueue.topic.name

错误记录报告器使用的 DLQ Topic 名称。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:非空时启用 Sink 错误记录报告路径;该 Topic 不能被 topics 消费或被 topics.regex 匹配。

errors.log.enable

是否启用 Kafka Connect 框架级错误日志。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:与 behavior.on.error=LOG 独立;只控制框架级错误日志。

最佳实践

首次接入结构化 Topic 并自动管理表

适用业务场景:首次把包含结构化 key/value 的 Kafka Topic 写入新的 ScyllaDB keyspace,希望 Connector 创建 keyspace 和表;后续需要新增普通列时,可以安排重启 Connector task,再由 Connector 检查并补充表结构。 配置示例
关键说明:记录 key 的字段会参与主键生成,value 提供其余列;实际 Kafka Topic 的句点和连字符会替换为下划线后决定表名。先确认 key 字段稳定且能唯一标识目标行,再让 Connector 初始化表结构。结构变化适合新增普通列;不要将它当作修改既有主键的机制。对于 ScyllaDB Sink Connector 1.1.9,同一 Task 生命周期内首次建立表后会命中 schema cache,后续记录中的新增普通列不会触发重新检查;需要让 Connector 检查并补充新增普通列时,必须重启处理该 Topic 分区的 Task,并在重启后再处理包含新列的记录。重启 Connector 本身不是这里的必要条件,关键是让负责处理记录的 Task 重新建立其缓存。

按 Topic 派生表并使用显式列映射

适用业务场景:多个 Kafka Topic 写入同一 keyspace 下的不同派生表,且需要只写入选定字段,或把 key、value、header 字段映射到明确的 ScyllaDB 列。 配置示例
关键说明:映射启用后未列出的字段不会写入;引用的 key/value 字段必须存在于记录 Schema,header 映射需要基本类型。动态配置键中的 <topic> 片段必须使用 Connector 对实际 Kafka Topic 规范化后的名称:句点和连字符会变成下划线,并据此查找 Topic 级配置。<keyspace><table> 片段仍需符合动态配置名称格式,不要把它们理解为会参与运行时 Topic 查找的独立路由值。

为不同 Topic 设置保留期和删除策略

适用业务场景:不同业务数据需要不同的保留时间,或只有部分 Topic 应将 tombstone 转换为按主键删除。 配置示例
关键说明:Topic 级 TTL 和删除开关会覆盖 Connector 级值。动态配置键中的 Topic 片段必须使用规范化后的 Topic 名称。TTL 作用于 ScyllaDB 的 INSERT 语句;tombstone 在删除开关启用且目标表存在时按记录 key 的全部主键字段执行 DELETE。两者都不改变 Kafka Offset 或事件时间语义。根据业务保留策略分别设置,并验证 tombstone 的上游产生方式。

监控

监控内容

关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在启用了相应错误处理时关注 DLQ 活动。

导入 Grafana 大盘

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

限制条件

  • topicstopics.regex 互斥,且必须至少配置一个非空订阅选择器;配置的 DLQ Topic 不能被订阅。
  • 实际 Kafka Topic 名称中的句点和连字符都会规范化为下划线;该规范化名称既用于派生表名,也用于查找 Topic 级动态配置。规范化后发生碰撞的 Topic 会指向同一个派生表。
  • 无 Schema 的 key 或 value 不用于生成或调整 DDL;使用无 Schema 数据前必须预先创建兼容的目标表。
  • 表结构管理只支持新增普通列;已有表的主键形状不能通过映射改变。
  • 在同一 task 生命周期内,schema cache 命中后不会重新检查记录中的新增普通列;要检查并补充新增普通列,必须重启 task。
  • Kafka Offset 提交与 ScyllaDB 业务写入不是同一事务;进程故障或重试可能再次处理记录,不能据此推断跨系统原子提交或重复免除。

常见问题

Connector 无法连接 ScyllaDB,如何排查?

检查 scylladb.contact.pointsscylladb.port 和网络访问;启用认证时确认用户名和密码同时设置,启用 TLS 时确认相关文件路径可读且证书配置匹配。修正连接配置后重新部署或重启 Connector,并观察 Task 是否恢复运行。

为什么配置了 topics 后仍无法启动?

检查是否同时配置了 topics.regex,或两个配置都为空。Sink 必须二选一使用非空的 topicstopics.regex;如果配置了 errors.deadletterqueue.topic.name,还要确保 DLQ Topic 不在订阅列表或正则匹配范围内。

为什么记录写入失败并提示 key 或 value 类型不支持?

确认 Worker Converter 输出的是 Struct 或 Map,而不是顶层基本类型、null key 或不支持的嵌套 Struct。检查记录 key 是否包含目标表的主键字段,value 字段类型是否与 ScyllaDB 列类型兼容;无 Schema Map 场景还要确认目标表已预创建。

tombstone 没有删除目标行怎么办?

确认 scylladb.deletes.enabled 或匹配 Topic 的 deletesEnabledtrue,并检查 tombstone 的 key 是否包含目标表的全部主键字段。若目标表不存在或 key 不完整,Connector 无法按主键构造有效删除语句;同时确认动态配置键中的 Topic 片段与规范化后的 Topic 名称匹配,并确认 keyspace/table 片段符合动态配置名称格式。