Cassandra作为一款广泛使用的分布式NoSQL数据库,在日志存储、物联网数据、用户行为记录等场景中表现出色。但当我们需要把这些数据实时同步到搜索引擎、消息队列或者缓存系统时,问题就来了:Cassandra没有像MySQL binlog那样通用的订阅机制,怎么做增量同步?全表轮询显然不现实,延迟高、负载大还容易漏数据。好在Cassandra原生提供了CDC(Change Data Capture,变更数据捕获)能力,通过它可以捕获表的写入和更新操作,实现准实时的数据流转。这篇文章就围绕Cassandra CDC的原理、配置和消费方式展开,帮你搭建一条可靠的数据同步链路。

Cassandra CDC的工作原理是什么
要理解CDC,得先从Cassandra的写入路径说起。Cassandra所有的写操作都会先顺序写入commitlog,然后再写入内存中的memtable,memtable满了之后刷盘形成SSTable。CDC正是利用了这个机制:当一张表开启了CDC属性后,其变更数据在被刷盘时会被单独提取出来,写入到CDC日志文件中。
具体来说,每个节点上都有一个名为cdc_raw的内部表,开启CDC的表的变更数据会以mutation的形式写到这个表里,最终落盘为CDC日志文件,存放在数据目录的cdc子目录下。这些文件遵循commitlog的段文件格式,后缀名为cdc.log。与普通commitlog不同的是,CDC日志不会在memtable刷盘后被清理,而是需要等待消费者确认处理完毕后手动删除,这正是数据不丢失的关键。
这里有一个容易踩的坑:CDC捕获的是写入操作,包括INSERT、UPDATE和DELETE,但它是节点本地的、物理层面的捕获,不是逻辑层面的行级变更流。也就是说,同一条记录被多次更新,CDC日志里会体现为多个mutation,消费端需要自己做合并或者按时间戳取最新值处理。
如何配置和启用CDC
启用CDC分两步,先在服务端开启功能,再在表级别声明。首先修改cassandra.yaml中的相关参数:
# cassandra.yaml 中的 CDC 相关配置 # 指定 CDC 日志可占用的最大磁盘空间,单位MB cdc_enabled: true cdc_free_space_in_mb: 4096 # 剩余空间低于该值时拒绝写入CDC表,单位MB cdc_free_space_check_interval_ms: 250
然后建表时加上WITH cdc = true选项,已存在的表也可以通过ALTER语句补加:
-- 新建开启CDC的表
CREATE TABLE user_events (
user_id UUID,
event_time TIMESTAMP,
event_type TEXT,
payload TEXT,
PRIMARY KEY (user_id, event_time)
) WITH cdc = true;
-- 对已有表启用CDC
ALTER TABLE user_events WITH cdc = true;
关于cdc_free_space_in_mb的设置需要特别小心。这个参数限制了CDC日志目录的最大空间,一旦达到上限,对应表的写入会被直接拒绝并抛出WriteFailureException。默认值4096MB在写入量大的场景下可能很快耗尽,建议根据下游消费能力评估:如果消费延迟通常在秒级,CDC日志积压量不会太大;如果下游偶尔停机维护,就要预留足够的缓冲空间。
同时建议配置好监控告警,密切关注CDCTotalSizeOnDisk这个指标(可以通过JMX或nodetool获取),当它超过阈值的80%时就该报警了,而不是等到写入被拒绝才发现问题。
如何消费CDC日志数据
CDC日志本质上是commitlog格式的文件,直接解析比较麻烦,官方Java驱动提供了CommitLogReader相关的工具类可以简化这个过程。基本思路是扫描CDC目录下的cdc.log文件,解析出mutation,再转换为具体的行数据。
下面是一个简化的消费框架示例,展示如何用Java驱动读取CDC文件并处理变更:
import org.apache.cassandra.db.commitlog.CommitLogReader;
public class CdcConsumer {
public void consume(File cdcFile) throws IOException {
CommitLogReader reader = new CommitLogReader();
// 读取CDC日志段文件,逐个mutation回调处理
reader.readAllFiles(cdcFile.getParentFile().toPath(),
(mutation, size) -> {
// 只处理目标表的变更
String tableName = CompactTableMapper.tableName(mutation);
if ("user_events".equals(tableName)) {
processMutation(mutation);
}
});
// 处理完成后,将cdc.log重命名为cdc.idx表示已消费
renameToProcessed(cdcFile);
}
private void processMutation(Mutation mutation) {
// 解析分区键和单元格,转换为事件对象后写入Kafka
mutation.getPartitionUpdates().forEach(update -> {
update.forEach(cell -> {
kafkaTemplate.send("cdc-topic", buildEvent(update, cell));
});
});
}
}
这里有个重要的约定:处理完一个CDC文件后,把它重命名为.cdc.idx后缀,表示该文件已被消费。Cassandra自身会定期清理这些已标记的文件(由cdc_segment_wait_ms参数控制,新版中默认是立即回收或10秒左右)。如果你的消费程序崩溃了,重启后需要跳过已处理的部分,因此建议在消费端维护一个处理进度的检查点,记录已完成的文件名和段ID。
除了自己写消费者,社区也有一些现成的方案,比如基于Debezium的Cassandra连接器(如Debezium Cassandra插件),它把CDC解析、位点管理、投递Kafka这些脏活都封装好了,适合不想重复造轮子的团队。不过要注意这类插件的版本兼容性问题,Debezium官方对Cassandra的支持长期处于incubating状态,选型前务必做好测试。
CDC与其他数据同步方案的对比
除了CDC,Cassandra生态里还有几种常见的数据同步思路,各有适用场景。最朴素的是基于时间戳的轮询,即定期查询WHERE event_time > 上次位置的数据。这种方式实现简单,但依赖业务表里有可靠的时间戳字段,且时钟偏移、删除操作无法感知(数据删了就查不到了),一般只适合对实时性要求不高的场景。
第二种是Cassandra的Trigger(触发器)机制,可以自定义一个实现了ITrigger接口的类,在写入时把变更转发到其他表或外部系统。Trigger的问题在于它是同步执行的,直接嵌入写入路径,一旦外部系统响应慢,会拖慢整个写入,而且实现类需要以JAR包形式部署到所有节点,运维成本不低。
第三种是双写方案,应用层在写Cassandra的同时写Kafka。这种方案把顺序性保证完全交给应用代码,一旦某一步失败就会出现数据不一致,而且侵入业务逻辑,后期难以维护,不推荐在核心链路使用。
对比下来,CDC的优势在于对写入路径零侵入、性能开销可控、能捕获包括DELETE在内的所有变更。劣势则是部署上是节点级别的,消费端需要聚合多个节点的数据才能得到全局视图,且原生工具链相对简陋,需要一定的开发投入。
生产环境的实践建议
最后总结几条实战经验。第一,CDC日志的空间管理是重中之重,务必给CDC目录设置独立的监控,同时确保消费程序的吞吐能力跟得上写入峰值,必要时在消费端做批量拉取和异步投递,避免日志积压触发写拒绝。
第二,做好幂等设计。CDC事件可能因为消费重试而被重复投递,下游系统(尤其是缓存更新、搜索索引这类操作)要以分区键加时间戳作为幂等键,保证重复消费不会产生脏数据。
第三,处理好节点扩缩容的情况。加入新节点后,新节点上的CDC日志需要被纳入消费范围,消费程序最好支持动态发现节点列表,可以通过查询system.local和system.peers表来获取当前拓扑,或者干脆部署一个集中式的分发层来统一管理。
第四,在测试环境充分验证删除事件的表现。DELETE在CDC中体现为一个带有删除标记的mutation(可能是row marker级别的删除,也可能是cell级别的墓碑),消费端要能正确区分并转换为下游的删除指令,这一点在同步到Elasticsearch之类的搜索引擎时尤其容易出问题。
总的来说,Cassandra CDC虽然没有Kafka Connect那样开箱即用的体验,但它提供了可靠的数据捕获基础,配合合理的消费架构,完全能够支撑起一条稳定的数据同步管道。理解了commitlog与CDC日志的关系,再把空间管理和消费位点这两件事做扎实,就能避开绝大多数坑。