概述
DataStax Cassandra Sink Connector 从 Kafka Topic 消费记录,并将消息 Key、Value 或 Header 中的字段映射到 Apache Cassandra 或 DataStax Enterprise 的目标表。它位于 Kafka 与 Cassandra 之间,适合把业务事件、实体状态和日志类数据持续写入按 Keyspace 和 Table 组织的存储模型。 每个 Topic 可以配置一个或多个目标表。Connector 根据映射生成 CQL 写入语句,也可以使用自定义 Prepared Statement;同一条 Kafka 记录写入多个表时,各表写入彼此独立。目标 Keyspace、Table、列和主键需要预先规划,映射负责明确 Kafka 字段与 Cassandra 列之间的关系。前置条件
- 在 Cassandra 中预先创建目标 Keyspace 和 Table,并确保连接账号具有读取表元数据以及执行所需写入或删除操作的权限;映射中引用的列必须存在,且自动生成 CQL 时必须覆盖全部主键列。
授权许可
使用 DataStax Apache Kafka Connector License Terms。发布包中的许可证标识存在差异,使用或分发前请核对随包许可证文件并完成客户侧确认。快速开始
提前准备 Connect Cluster、Kafka、Cassandra Keyspace 和目标表,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector。 以下示例订阅orders Topic,将消息 Key 写入 order_id,并将消息 Value 中的 customer_id 和 status 写入 app.orders 表。目标表应以 order_id 作为主键,并包含这些列。
<cassandra-host> 替换为 Cassandra 联系节点,将 <local-dc> 替换为该节点所属的数据中心名称。如果 Topic、Keyspace、Table 或字段名不同,请同步修改 topics、动态配置键和 mapping;启动后发送包含稳定 Key 以及 customer_id、status 字段的记录,并检查 app.orders 中对应行。
配置
Kafka Connect 与订阅
connector.class
指定要加载的 Connector 实现类。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:使用
com.datastax.oss.kafka.sink.CassandraSinkConnector。旧类名com.datastax.kafkaconnector.DseSinkConnector仅作为已弃用的兼容别名。 - 必填:是
tasks.max
设置允许启动的最大 Task 数量。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:至少为
1。实际 Task 数不会超过可分配的 Topic Partition 数量;增加 Task 不保证吞吐线性增长。
topics
显式列出要消费的 Kafka Topic。
- 类型:
list - 默认值:空列表
- 重要级别:高
- 有效值 / 注意事项:以逗号分隔。必须在
topics和topics.regex中仅配置一个;每个显式 Topic 都必须至少有一个对应的topic.<topic>.<keyspace>.<table>.*表配置。
topics.regex
使用 Java 正则表达式订阅 Kafka Topic。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:必须在
topics和topics.regex中仅配置一个。正则新匹配到的 Topic 仍需有名称完全对应的topic.<topic>.*配置,否则记录到达时会映射失败。
Cassandra 连接
contactPoints
设置 Cassandra 初始联系节点。
- 类型:
list - 默认值:空列表
- 重要级别:高
- 有效值 / 注意事项:填写逗号分隔的 IP 地址或域名;所有节点使用
port指定的同一端口。非空时还必须配置本地数据中心,且不能与cloud.secureConnectBundle同时使用。需要为不同节点指定不同端口时,改用 Java Driver 原生联系点配置。
port
设置 contactPoints 使用的 Cassandra 原生传输端口。
- 类型:
int - 默认值:
9042 - 重要级别:高
- 有效值 / 注意事项:至少为
1;使用安全连接包时忽略此项。
loadBalancing.localDc
设置联系节点所属的本地数据中心。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:
contactPoints非空时必须配置;使用安全连接包时必须留空。若同时配置替代项,本项优先。 - 已弃用:是
- 替代项:
datastax-java-driver.basic.load-balancing-policy.local-datacenter
cloud.secureConnectBundle
设置 DataStax Cloud 安全连接包路径。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:非空时启用云连接模式,不能同时配置
contactPoints、loadBalancing.localDc或任何ssl.*项。云模式会把ANY、ONE和LOCAL_ONE写一致性级别调整为LOCAL_QUORUM。若同时配置替代项,本项优先。 - 已弃用:是
- 替代项:
datastax-java-driver.basic.cloud.secure-connect-bundle
认证
auth.provider
选择 Cassandra 认证提供者。
- 类型:
string - 默认值:
None - 重要级别:高
- 有效值 / 注意事项:区分大小写,可选
None、PLAIN或GSSAPI。配置用户名或密码时,即使本项为None,实际也会使用PLAIN。
auth.username
设置 PLAIN 认证用户名。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:用于
PLAIN认证;配置非空密码时用户名也必须非空。
auth.password
设置 PLAIN 认证密码。
- 类型:
password - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:用于
PLAIN认证。通过安全的配置提供方式注入,不要在日志或共享文件中暴露明文密码。
auth.gssapi.keyTab
设置 Kerberos keytab 文件路径。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:非空路径必须指向 Worker 可读取的普通文件。使用 keytab 且未配置 principal 时,可从 keytab 的第一个 principal 推断;不使用 keytab 时由 Driver 使用票据缓存。
auth.gssapi.principal
设置 Kerberos principal。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:仅用于
GSSAPI。省略该项与显式配置空字符串的行为不同;只有省略时才能从 keytab 推断。
auth.gssapi.service
设置 GSSAPI SASL 服务名。
- 类型:
string - 默认值:
dse - 重要级别:高
- 有效值 / 注意事项:使用
GSSAPI时必须为非空字符串,并与 Cassandra 服务端配置一致。
TLS
ssl.provider
选择 TLS 实现提供者。
- 类型:
string - 默认值:
None - 重要级别:高
- 有效值 / 注意事项:区分大小写,可选
None、JDK或OpenSSL。使用安全连接包时不能配置任何ssl.*项。
ssl.cipherSuites
设置允许使用的 TLS 密码套件。
- 类型:
list - 默认值:空列表
- 重要级别:高
- 有效值 / 注意事项:以逗号分隔;空列表表示由所选 TLS 提供者使用默认值。不能与安全连接包同时配置。
ssl.hostnameValidation
控制是否校验 Cassandra 节点主机名。
- 类型:
boolean - 默认值:
true - 重要级别:高
- 有效值 / 注意事项:为
true时,JDK SSL 会解析联系节点地址并启用主机名校验。不能与安全连接包同时配置。
ssl.keystore.password
设置 JDK SSL keystore 密码。
- 类型:
password - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:用于访问 JDK SSL keystore。通过安全方式注入,且不能与安全连接包同时配置。
ssl.keystore.path
设置 JDK SSL keystore 路径。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:非空路径必须指向 Worker 可读取的普通文件。不能与安全连接包同时配置。
ssl.openssl.keyCertChain
设置 OpenSSL 客户端证书链路径。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:非空路径必须指向 Worker 可读取的普通文件。使用
OpenSSL时,必须与ssl.openssl.privateKey同时配置或同时省略;不能与安全连接包同时配置。
ssl.openssl.privateKey
设置 OpenSSL 客户端私钥路径。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:非空路径必须指向 Worker 可读取的普通文件。使用
OpenSSL时,必须与ssl.openssl.keyCertChain同时配置或同时省略;保护私钥文件权限,且不能与安全连接包同时配置。
ssl.truststore.password
设置 TLS truststore 密码。
- 类型:
password - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:用于 JDK 或 OpenSSL 初始化时加载 truststore。通过安全方式注入,且不能与安全连接包同时配置。
ssl.truststore.path
设置 TLS truststore 路径。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:非空路径必须指向 Worker 可读取的普通文件;OpenSSL 模式按 JKS 加载该文件。不能与安全连接包同时配置。
写入吞吐与请求
maxConcurrentRequests
限制同一 Worker JVM 中该 Connector 实例的在途 Cassandra 请求数。
- 类型:
int - 默认值:
500 - 重要级别:高
- 有效值 / 注意事项:至少为
1。同名 Connector 的多个 Task 在同一 Worker JVM 中共享此上限;跨 Worker 不共享。
maxNumberOfRecordsInBatch
设置一个 Cassandra 批请求中的最大记录数。
- 类型:
int - 默认值:
32 - 重要级别:高
- 有效值 / 注意事项:至少为
1。Connector 按 Topic、目标表和 Cassandra Routing Key 分桶;多条语句使用 UNLOGGED batch,单条语句直接异步执行。
compression
设置 Cassandra 协议压缩算法。
- 类型:
string - 默认值:
None - 重要级别:高
- 有效值 / 注意事项:不区分大小写,可选
none、snappy或lz4。若同时配置替代项,本项优先。 - 已弃用:是
- 替代项:
datastax-java-driver.advanced.protocol.compression
queryExecutionTimeout
设置 CQL 请求超时秒数。
- 类型:
int - 默认值:
30 - 重要级别:高
- 有效值 / 注意事项:至少为
1,值按秒传递给 Java Driver。若同时配置替代项,本项优先。 - 已弃用:是
- 替代项:
datastax-java-driver.basic.request.timeout
connectionPoolLocalSize
设置每个本地 Cassandra 节点的连接池大小。
- 类型:
int - 默认值:
4 - 重要级别:高
- 有效值 / 注意事项:至少为
1。若同时配置替代项,本项优先。 - 已弃用:是
- 替代项:
datastax-java-driver.advanced.connection.pool.local.size
错误处理与指标
ignoreErrors
控制 Connector 是否忽略记录映射或 Cassandra Driver 写入错误。
- 类型:
string - 默认值:
None - 重要级别:高
- 有效值 / 注意事项:不区分大小写,可选
None、Driver或All。Driver仅忽略数据库 Driver 写入失败,All还忽略映射失败;忽略错误会允许 Offset 越过失败记录,Connector 不会自动补写这些数据。该策略独立于 Kafka Connect 的errors.tolerance。
jmx
控制是否启用默认 Java Driver JMX 会话指标。
- 类型:
boolean - 默认值:
true - 重要级别:高
- 有效值 / 注意事项:为
true且未显式配置 Driver 指标列表时,启用cql-requests和cql-client-timeouts,其中cql-requests的默认采样间隔为 30 秒。
metricsHighestLatency
设置 CQL 请求指标的最高延迟量程秒数。
- 类型:
int - 默认值:
35 - 重要级别:高
- 有效值 / 注意事项:至少为
1,通常应大于请求超时;该大小关系不会被强制校验。若同时配置替代项,本项优先。 - 已弃用:是
- 替代项:
datastax-java-driver.advanced.metrics.session.cql-requests.highest-latency
Topic 数据转换
topic.<topic>.codec.locale
设置指定 Topic 的文本转换区域。
- 类型:
string - 默认值:
en_US - 重要级别:高
- 有效值 / 注意事项:
<topic>只能包含字母、数字、句点、下划线和连字符;区域值必须可由 DSBulk Codec 解析。
topic.<topic>.codec.timeZone
设置指定 Topic 的时间转换时区。
- 类型:
string - 默认值:
UTC - 重要级别:高
- 有效值 / 注意事项:必须是
java.time.ZoneId可识别的时区,例如Asia/Shanghai。
topic.<topic>.codec.timestamp
设置字符串转换为 CQL timestamp 时使用的格式。
- 类型:
string - 默认值:
CQL_TIMESTAMP - 重要级别:高
- 有效值 / 注意事项:可使用 DSBulk Codec 支持的时间模式、
DateTimeFormatter常量或CQL_TIMESTAMP。
topic.<topic>.codec.date
设置字符串转换为 CQL date 时使用的格式。
- 类型:
string - 默认值:
ISO_LOCAL_DATE - 重要级别:高
- 有效值 / 注意事项:可使用 DSBulk Codec 支持的日期模式或
DateTimeFormatter常量。
topic.<topic>.codec.time
设置字符串转换为 CQL time 时使用的格式。
- 类型:
string - 默认值:
ISO_LOCAL_TIME - 重要级别:高
- 有效值 / 注意事项:可使用 DSBulk Codec 支持的时间模式或
DateTimeFormatter常量。
topic.<topic>.codec.unit
设置纯数字时间输入的时间单位。
- 类型:
string - 默认值:
MILLISECONDS - 重要级别:高
- 有效值 / 注意事项:必须精确使用
NANOSECONDS、MICROSECONDS、MILLISECONDS、SECONDS、MINUTES、HOURS或DAYS。
表映射与写入语义
topic.<topic>.<keyspace>.<table>.mapping
定义 Kafka 字段到 Cassandra 列或自定义查询绑定变量的映射。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:每个目标表都必须提供非空映射,格式为逗号分隔的
列=字段表达式。字段表达式支持key、value、key.*、value.*、header.*、now(),以及保留伪列__ttl和__timestamp。自动生成 CQL 时,普通列必须存在且所有主键列都必须映射。 - 必填:是
topic.<topic>.<keyspace>.<table>.deletesEnabled
控制是否根据完整映射记录生成整行删除。
- 类型:
boolean - 默认值:
true - 重要级别:高
- 有效值 / 注意事项:只有未使用自定义查询且映射覆盖目标表全部列时,主键外所有映射值均为
null的记录才会删除整行。配置自定义query时必须设为false。
topic.<topic>.<keyspace>.<table>.consistencyLevel
设置目标表写入的一致性级别。
- 类型:
string - 默认值:
LOCAL_ONE - 重要级别:高
- 有效值 / 注意事项:使用 DataStax Driver 支持的一致性级别名称,不区分大小写。云模式会把
ANY、ONE和LOCAL_ONE调整为LOCAL_QUORUM。
topic.<topic>.<keyspace>.<table>.ttl
设置目标表写入的固定 TTL。
- 类型:
int - 默认值:
-1 - 重要级别:高
- 有效值 / 注意事项:至少为
-1;-1表示禁用固定 TTL。值按ttlTimeUnit转换为秒;映射中同时提供__ttl时动态值优先。Counter 表不能使用 TTL,自定义查询自行决定 TTL 语义。
topic.<topic>.<keyspace>.<table>.nullToUnset
控制非主键 null 是否按 Cassandra UNSET 处理。
- 类型:
boolean - 默认值:
true - 重要级别:高
- 有效值 / 注意事项:为
true时,非主键null不覆盖现有列值;为false时显式绑定null。主键为null始终失败。
topic.<topic>.<keyspace>.<table>.ttlTimeUnit
设置固定 TTL 和映射 __ttl 值的输入单位。
- 类型:
string - 默认值:
SECONDS - 重要级别:高
- 有效值 / 注意事项:必须精确使用
NANOSECONDS、MICROSECONDS、MILLISECONDS、SECONDS、MINUTES、HOURS或DAYS。
topic.<topic>.<keyspace>.<table>.timestampTimeUnit
设置映射 __timestamp 值的输入单位。
- 类型:
string - 默认值:
MICROSECONDS - 重要级别:高
- 有效值 / 注意事项:必须使用
TimeUnit枚举值。未映射__timestamp且未使用自定义查询时,Connector 会把 Kafka 记录时间戳从毫秒转换为微秒作为 CQL 写时间戳。
topic.<topic>.<keyspace>.<table>.query
为目标表设置自定义 Prepared Statement CQL。
- 类型:
string - 默认值:
null - 重要级别:高
- 有效值 / 注意事项:非
null值会替代自动生成的 INSERT 或 Counter UPDATE,并要求deletesEnabled=false。绑定变量必须由mapping提供;Connector 不再自动校验表列、主键、TTL、时间戳和删除语义。空字符串也会被视为已配置并在数据库 prepare 时失败。
Java Driver 透传
datastax-java-driver.<driver-path>
将 DataStax Java Driver 原生配置传递给 Driver。
- 类型:
string - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:使用
datastax-java-driver.<driver-path>形式,路径、取值、默认值和校验规则由随 Connector 提供的 Java Driver 决定。联系点、刷新 Keyspace、节点指标、会话指标和 TLS 密码套件这五类列表路径会按逗号拆分;旧 Connector 别名与对应 Driver 路径同时配置时,旧别名优先。
最佳实践
同步业务删除到 Cassandra
适用业务场景:基础写入已稳定运行,上游会用“保留主键、其余业务字段均为null”的记录表示实体删除,希望 Cassandra 中对应整行随之删除。目标表列结构稳定,并且映射可以覆盖表中的全部列。
配置示例:在快速开始配置中把删除开关改为 true。下面列出完整 Connector 配置,便于独立验证删除前的写入和后续删除结果。
customer_id 和 status 等所有非主键映射值均为 null 时,Connector 按 order_id 删除整行。若只映射部分列,记录仍走写入路径。自定义 query 与自动删除互斥,不要同时启用。
按 Topic Partition 扩展写入吞吐
适用业务场景:单 Task 已能正确写入,但 Kafka 消费积压持续增长,orders Topic 有多个 Partition,Cassandra 集群和 Worker 仍有可用容量,希望逐步提高并行消费和在途请求量。
配置示例:以下配置在快速开始的安全写入基线上将 Task 上限提高到 4,并使用明确的请求与批次起始值。上线前先在压测环境验证,再根据实际延迟和 Cassandra 负载逐项调整。
tasks.max=4 只是 Task 上限,实际并行度不会超过 Topic Partition 数量。maxConcurrentRequests 是同一 Worker JVM 内同名 Connector 各 Task 的共享上限;增加 Worker 还会增加独立 Session 和总并发。批次按目标表和 Cassandra Routing Key 分桶,64 不是通用最优值。调整前后应比较 Kafka Lag、写入吞吐、请求延迟、超时和 Cassandra 负载,且不要依赖跨 Partition、Task 或批次的全局顺序。
监控
监控内容
关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在启用了相应错误处理时关注 DLQ 活动。导入 Grafana 大盘
确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。限制条件
- Connector 不自动创建 Keyspace、Table 或列,也不会自动适配不兼容的表结构变化;修改目标 Schema 或映射后应重启 Task,使其重新读取元数据并准备 CQL。
- Cassandra 写入与 Kafka Offset 提交不是原子事务;故障恢复可能重放已成功写入但尚未提交 Offset 的记录,不能据此获得 exactly-once 或无重复保证。
- 同一记录写入多个目标表时不提供跨表事务;部分表写入成功后发生失败,重放可能再次写入已成功的表。
- Counter、
now()、TTL、自定义 CQL 和删除具有各自的重放副作用,不能假定普通 Upsert 的幂等特征适用于这些写入方式。
常见问题
Task 启动时提示映射或目标列无效
目标 Keyspace、Table 或列不存在,映射遗漏主键列,或者动态配置键中的 Topic、Keyspace、Table 与实际资源不一致,都可能导致启动失败。先核对目标表 Schema,再检查topic.<topic>.<keyspace>.<table>.mapping 中每个目标列和字段表达式;自动生成 CQL 时确保所有主键列均已映射,修改后重启 Connector。
Connector 已运行但部分记录没有写入 Cassandra
记录缺少映射要求的字段、主键值为null、类型无法转换或 Cassandra 请求失败时,写入可能失败;启用 ignoreErrors=Driver 或 All 还可能允许 Offset 越过失败记录。检查 Task 日志、错误计数和消费 Offset,确认映射字段与消息结构一致,并优先使用 ignoreErrors=None 排查和修复根因;被忽略的记录需要按业务流程单独补偿。