Skip to main content

概述

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_idstatus 写入 app.orders 表。目标表应以 order_id 作为主键,并包含这些列。
<cassandra-host> 替换为 Cassandra 联系节点,将 <local-dc> 替换为该节点所属的数据中心名称。如果 Topic、Keyspace、Table 或字段名不同,请同步修改 topics、动态配置键和 mapping;启动后发送包含稳定 Key 以及 customer_idstatus 字段的记录,并检查 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
  • 默认值:空列表
  • 重要级别:高
  • 有效值 / 注意事项:以逗号分隔。必须在 topicstopics.regex 中仅配置一个;每个显式 Topic 都必须至少有一个对应的 topic.<topic>.<keyspace>.<table>.* 表配置。

topics.regex

使用 Java 正则表达式订阅 Kafka Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:必须在 topicstopics.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
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:非空时启用云连接模式,不能同时配置 contactPointsloadBalancing.localDc 或任何 ssl.* 项。云模式会把 ANYONELOCAL_ONE 写一致性级别调整为 LOCAL_QUORUM。若同时配置替代项,本项优先。
  • 已弃用:是
  • 替代项datastax-java-driver.basic.cloud.secure-connect-bundle

认证

auth.provider

选择 Cassandra 认证提供者。
  • 类型string
  • 默认值None
  • 重要级别:高
  • 有效值 / 注意事项:区分大小写,可选 NonePLAINGSSAPI。配置用户名或密码时,即使本项为 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
  • 重要级别:高
  • 有效值 / 注意事项:区分大小写,可选 NoneJDKOpenSSL。使用安全连接包时不能配置任何 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
  • 重要级别:高
  • 有效值 / 注意事项:不区分大小写,可选 nonesnappylz4。若同时配置替代项,本项优先。
  • 已弃用:是
  • 替代项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
  • 重要级别:高
  • 有效值 / 注意事项:不区分大小写,可选 NoneDriverAllDriver 仅忽略数据库 Driver 写入失败,All 还忽略映射失败;忽略错误会允许 Offset 越过失败记录,Connector 不会自动补写这些数据。该策略独立于 Kafka Connect 的 errors.tolerance

jmx

控制是否启用默认 Java Driver JMX 会话指标。
  • 类型boolean
  • 默认值true
  • 重要级别:高
  • 有效值 / 注意事项:为 true 且未显式配置 Driver 指标列表时,启用 cql-requestscql-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
  • 重要级别:高
  • 有效值 / 注意事项:必须精确使用 NANOSECONDSMICROSECONDSMILLISECONDSSECONDSMINUTESHOURSDAYS

表映射与写入语义

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

定义 Kafka 字段到 Cassandra 列或自定义查询绑定变量的映射。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:每个目标表都必须提供非空映射,格式为逗号分隔的 列=字段表达式。字段表达式支持 keyvaluekey.*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 支持的一致性级别名称,不区分大小写。云模式会把 ANYONELOCAL_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
  • 重要级别:高
  • 有效值 / 注意事项:必须精确使用 NANOSECONDSMICROSECONDSMILLISECONDSSECONDSMINUTESHOURSDAYS

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 配置,便于独立验证删除前的写入和后续删除结果。
关键说明:自动删除不要求 Kafka Value 本身是墓碑,但映射必须覆盖目标表全部列;当 customer_idstatus 等所有非主键映射值均为 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=DriverAll 还可能允许 Offset 越过失败记录。检查 Task 日志、错误计数和消费 Offset,确认映射字段与消息结构一致,并优先使用 ignoreErrors=None 排查和修复根因;被忽略的记录需要按业务流程单独补偿。

修改 Cassandra 表结构后写入开始失败

Connector 在 Task 启动时读取表元数据、校验映射并准备 CQL,运行中的兼容性取决于已缓存的映射和 Prepared Statement。完成 Schema 变更后,先确认映射仍覆盖有效列和主键,再重启相关 Task 触发重新校验;对于删除列、修改主键或改变类型等不兼容变更,应先规划迁移和消息兼容方案。