
大数据这东西越做到后面越会发现一个扎心的事实没有一种存储能同时把写入、检索、分析全都干漂亮。接触过 Cassandra 和 Elasticsearch 的人应该都有过这种纠结——业务量一上来数据落库要扛住千万级并发写产品方又天天追着要模糊搜索、聚合统计单靠任何一边都像让一个人既当仓库管理员又当图书管理员最后两边都做不好。这篇文章我就把 Cassandra 和 Elasticsearch 整合这件事掰开揉碎讲清楚不整虚的从架构选型到环境部署从数据管道到一致性对账全程用我实际趟过的经验说话适合正在做大数据相关毕设、或者生产系统里准备引入这套组合的团队参考。1. 为什么非要把 Cassandra 和 Elasticsearch 绑在一起1.1 Cassandra 的海量写入能力与查询软肋Cassandra 最抓人的一点就是它的写入路径几乎是为吞吐量设计的。底层用 LSM-Tree写入先落 commit log再进 memtable批量刷盘成 SSTable整个过程没有随机 IO也没有传统数据库行锁的竞争所以水平扩节点之后写性能基本能线性增长。我在压测环境里用 6 个节点跑过单条消息几百字节的日志数据单集群能稳定吃掉每秒二三十万条写入这个吞吐量在开源存储里非常能打。但问题也出在存储模型上。Cassandra 的 CQL 查询本质上是围绕分区键设计的只要你的查询条件里带了主键或分区键它能在毫秒级定位数据可一旦查询条件不走主键比如按订单号模糊查或者按用户昵称搜这种需求Cassandra 会直接告诉你ALLOW FILTERING意思是它要把全表数据拉到内存里逐条过滤数据量到千万级别就是灾难现场。还有一点它不支持真正的全文索引、不支持 join连LIKE %xx%这种基础需求都做不了。拿我经历过的一个订单系统来说单日新增订单 3000 万条订单状态要查、要按用户维度汇总还要支持商品名称模糊查询。纯 Cassandra 架构下订单主键查明细没问题但业务方一句帮我查一下商品名里带华为的订单就直接把集群查挂了。那一刻你就明白Cassandra 天生不是干检索的。1.2 Elasticsearch 的检索强项与写入瓶颈Elasticsearch 的本事不用多夸倒排索引、分词器、聚合分析简直就是为检索场景生的。它的原理类似于给数据建立目录索引查询时不用扫全表而是直接通过词项找文档列表所以海量数据里做模糊查询、组合条件筛选、聚合统计都是毫秒级响应。但 ES 的软肋同样明显写入链路太长。一条文档进来要写 translog要进内存 buffer要定期 refresh 生成新的 segment后台还有 merge 线程不停地合并小 segment。高峰期并发一高CPU 先被打满紧接着就是RejectedExecutionException集群拒绝写入。我见过很多团队把 Cassandra 换成 ES 想解决检索问题结果数据量一上来每天凌晨的批量任务直接把 ES 写崩溃最后还得降级成半夜低峰同步治标不治本。1.3 写出读搜的组合逻辑Cassandra 负责扛住所有的写入洪峰同时作为数据的最终权威存储ES 则通过异步管道拿到同一份数据专职承担搜索和聚合。查询请求只打到 ES写入请求只打到 Cassandra两边各干各的谁也不拖累谁。这个架构在日志检索、用户行为分析、订单流水查询以及 IoT 设备上报这类写入量大、检索需求复杂的场景里基本属于标配方案。核心思想可以用一句话概括写入以 Cassandra 为准查询以 ES 为准数据通过管道保持最终一致。2. 整合架构的三种主流玩法2.1 双写模式最快落地但别上来就梭哈最简单的整合方式就是在业务代码里同时写两个存储先插入 Cassandra成功后调用 ES 接口写入文档。这种方式适合项目初期或者对实时性要求特别高的场景代码直观、链路短、出问题好排查。但双写有几个躲不开的坑。最典型的是事务问题Cassandra 写成功ES 写超时了这时候两边数据就对不上了。有人会搞先写 ES 再写 Cassandra那反过来也一样有窗口。另外业务高峰期如果同步写入 ES等于又把压力传导给了 ES这就是双写方案和给 ES 减负的初衷互相矛盾的地方。我的建议是双写只适合以下两种情况数据量还没上来ES 写入压力不大且业务能容忍少量数据不一致作为临时方案跑通流程后续仍然要迁移到异步管道。如果决定用双写务必要有一个补偿机制。比较务实的做法是Cassandra 写入成功之后把一条待同步事件放进本地消息表或者 MQ如果 ES 写入失败后续由消费者从消息表里捞出来重试。这样主链路虽然看起来是双写底层其实已经具备了异步重放的兜底能力。2.2 CDC 异步管道生产环境的长期正确解我在线上项目里最终落地的方案是基于 CDC 的异步管道。完整链路长这样Cassandra 的 commit log 里记录了每一次数据变更打开 CDC 后变更事件会被写进专门的 cdc_raw 目录然后由捕获组件比较常用的是 Debezium 的 Cassandra Connector读取这些变更文件把新增、修改、删除操作转成事件消息发到 Kafka下游再起一个消费服务从 Kafka 拉消息经过处理后批量写入 Elasticsearch。你可能会问为什么中间非要插一个 Kafka直接从 Cassandra 读出来写 ES 不行吗原因有三个。第一是削峰填谷Cassandra 的写入洪峰ES 不一定扛得住Kafka 作为缓冲区可以把峰值流量摊平ES 按自己的消费能力慢慢消化。第二是失败重试消费端写 ES 失败时消息还能留在 Kafka 里重新消费不会丢数据。第三是解耦如果以后除了 ES 还需要同步到别的存储比如 Redis 缓存、ClickHouse只需新增一个消费者Cassandra 那一侧完全不用动。这个方案的具体配置和代码我在第 4 节详细展开。这里想强调一个容易忽视的地方Cassandra 的 CDC 目录是会被写满的。因为变更日志文件不会自动清理需要消费端在读完文件后及时把.cdc文件移除否则磁盘空间迟早被打爆而且 Cassandra 在 CDC 目录空间不足时会直接拒绝新的变更写入影响面非常大。2.3 迁移期的离线全量 实时增量组合如果你的 Cassandra 里已经积累了几个月甚至几年的历史数据直接搭好管道等增量是远远不够的ES 里必须先把存量数据灌进去否则用户搜索历史数据什么都搜不到。这时候就需要先全量、再增量两步走。全量阶段用 Spark 或者写一个独立的读任务按主键范围分批从 Cassandra 扫描数据然后通过 ES 的_bulk接口批量写入。注意全量期间要暂时关闭 ES 索引的刷新频率和副本数否则写入速度会被拖垮这个我后面在调优部分会细说。全量完成之后再启动 CDC 增量管道并且要保证增量的起点不早于全量扫描的开始时间否则重叠窗口内可能出现漏数据。实际操作中我会在全量任务和增量任务之间做个简单的版本对齐给表加一个last_updated时间戳字段全量启动时记录当前时间T0增量消费端从 Kafka 里取事件时只处理时间T0之后的事件这样既不会重复处理全量已经灌进去的数据也不会漏掉全量期间新写入的数据。3. 环境搭建最容易翻车的几个地方3.1 版本兼容性才是第一关不少人在整合项目里遇到的第一个拦路虎不是代码写不对而是组件版本互相打架。Cassandra 和 Elasticsearch 背后都有各自的 JDK 要求版本不匹配时经常出现服务启动一半直接退出甚至启动成功但内部行为异常的情况。我把常用版本对应的 JDK 要求整理成了一张表方便你对照组件推荐版本JDK 要求备注Cassandra 3.11.x生产仍有人用JDK 8版本较老新项目不建议Cassandra 4.0.x稳定主流JDK 8 或 11生产常见Cassandra 4.1.x较新JDK 8/11/17注意部分 JVM 参数在不同 JDK 下不兼容Elasticsearch 7.x常见内置 JDK需 JDK 11使用内置 JDK 时不要额外指定 JDK 17Elasticsearch 8.x / 9.x当前主流内置 JDK 17对系统资源要求更高值得提醒的是ES 虽然内置了 JDK但它会自动检测JAVA_HOME环境变量如果你的机器上装的是 JDK 17 而启动 ES 7.x很可能直接报Unsupported Java version。解决方法是把JAVA_HOME指向 ES 自带的 JDK 目录或者在启动脚本里显式指定ES_JAVA_HOME。Cassandra 那边如果用了 JDK 17 并把某个版本里默认开启的-XX:UseConcMarkSweepGC参数照搬过来会直接报Unrecognized VM option折腾半天才发现是 GC 参数不兼容。3.2 内存和系统参数ES 的堆内存设置有个著名的黄金法则堆大小不超过物理内存的一半且不要超过 31GB。超过 31GB 后 JVM 的压缩对象指针会失效内存反而浪费。一般建议给 ES 堆设置 16GB-31GB剩下的内存留给操作系统文件缓存因为 ES 的倒排索引和文档缓存大量依赖 OS cache这个比堆内存更重要。Cassandra 的内存规划则不同它的堆内存主要用于 memtable 和缓存通常设到 8GB 就够用但一定要给 JVM off-heap 和页缓存留足空间。很多 Cassandra 集群性能差不是 CPU 不行而是内存没分配好。Linux 系统下记得关闭 swap执行swapoff -a否则内存抖动时 JVM 表现会很差。同时把文件句柄数调大ES 和 Cassandra 都是文件密集型应用默认的 1024 肯定不够建议设置ulimit -n 65535以上。Windows 本地开发时这些问题不太明显但 ES 启动时如果看到max virtual memory areas vm.max_map_count [65530] is too low的提示这个在 Windows 上可以先忽略等真正部署到 Linux 时再执行sysctl -w vm.max_map_count262144解决。3.3 Windows 本地开发环境的启动细节我知道很多人是在 Windows 上做开发和毕设的所以多说几句。Elasticsearch 在 Windows 下的启动方式很简单下载 zip 包解压后进 bin 目录双击elasticsearch.bat就行。常见的坑有这么几个高版本 ES 启动后访问http://localhost:9200会要求安全认证首次启动会在控制台打印一个临时密码和证书指纹如果不想折腾可以在elasticsearch.yml里把xpack.security.enabled设为false关掉认证仅限本地开发。9200 是 HTTP 端口9300 是节点间传输端口如果启动报端口被占用用netstat -ano | findstr 9200找到占用进程并处理。ES 默认不允许以管理员身份直接跑在 root 下在 Windows 上虽然没有 root 概念但某些杀毒软件会对 ES 的数据目录做实时扫描严重影响启动和写入性能建议把 ES 和 Cassandra 的数据目录加到杀毒白名单里。Cassandra 在 Windows 上通过bin\cassandra.bat -f前台启动第一次启动要等一会儿看到Starting listening for CQL clients on /127.0.0.1:9042就说明起来了。注意 Cassandra 的默认数据目录在系统盘大量测试数据会占用不少空间最好提前改cassandra.yaml里的data_file_directories。3.4 最小集群规模建议有人图省事想用单节点 Cassandra 单节点 ES 先跑起来。我只能说这适合本地验证功能千万别把这种认知带到生产。Cassandra 单节点意味着没有副本节点一挂数据全丢ES 单节点意味着没有副本分片重启或宕机期间索引不可用。生产环境最少也要 Cassandra 3 节点起步副本因子设为 3本地机架写 2 份、跨机架写 1 份ES 同样 3 节点起步每个索引的主分片设为 3、副本设为 1。至于是否把 Cassandra 和 ES 部署在同一批机器上小规模集群为了节省成本可以混部但前提是给两边的 JVM 内存和 CPU 做好隔离比如用 cgroup 限制。数据量上来之后强烈建议分离部署因为 Cassandra 是 IO 密集型ES 是 CPU 和内存密集型混部容易互相干扰定位问题时也会更麻烦。4. 数据管道从零实现双写与 CDC 的完整例子4.1 双写方案的代码演示先看双写怎么实现。我习惯用 Python 做快速验证生产用 Java但核心逻辑是一样的。这里给你一套可以直接跑通思路的代码。from cassandra.cluster import Cluster from elasticsearch import Elasticsearch, helpers import uuid # 连接 Cassandra cluster Cluster([127.0.0.1], port9042) session cluster.connect(sales) # 连接 Elasticsearch es Elasticsearch([http://127.0.0.1:9200]) def write_order(order_id, user_id, goods_name, amount): # 1. 写入 Cassandra session.execute( INSERT INTO orders (order_id, user_id, goods_name, amount, created_at) VALUES (%s, %s, %s, %s, toTimestamp(now())), (order_id, user_id, goods_name, amount) ) # 2. 写入 Elasticsearch doc { order_id: str(order_id), user_id: str(user_id), goods_name: goods_name, amount: amount } es.index(indexorders, idstr(order_id), documentdoc) # 调用 write_order(uuid.uuid4(), u_10001, 华为手机 Mate 60, 6999.00)这段代码的问题在于第 2 步如果抛异常函数直接报错业务上不能说订单没创建成功但 ES 里又没有数据。前文说的补偿机制在这里就能用上把待同步记录先存进一张sync_task表后台线程定时扫描这张表把失败的数据重新灌给 ES。CREATE TABLE sales.sync_task ( task_id uuid PRIMARY KEY, order_id uuid, retry_count int, last_error text, status text );这样双写方案就从写完拉倒升级成了写失败可重试可靠性大幅提升。4.2 CDC 管道的关键配置与消费端逻辑生产环境我更推荐 CDC 方案。首先在 Cassandra 里开启 CDC修改cassandra.yamlcdc_enabled: true cdc_raw_directory: /var/lib/cassandra/cdc_raw cdc_total_space_in_mb: 4096同时需要为具体表开启 CDC。执行 CQLALTER TABLE sales.orders WITH cdc true;之后凡是 orders 表有数据变更Cassandra 会把变更日志写到 cdc_raw 目录。这里我用的捕获工具是 Debezium 的 Cassandra Connector它的配置大致长这样{ name: cassandra-cdc-connector, config: { connector.class: io.debezium.connector.cassandra.CassandraConnector, tasks.max: 1, topic.prefix: cdc-events, cassandra.hosts: node1,node2,node3, cassandra.cdc.dir: /var/lib/cassandra/cdc_raw, snapshot.mode: initial, table.include.list: sales.orders } }连接器启动后变更事件会以 JSON 格式发到 Kafka 的cdc-events.sales.orders主题。消费端拿到事件后需要做三件事解析事件里的 before/after 结构、根据操作类型决定写入还是删除、最后批量写入 ES。from kafka import KafkaConsumer from elasticsearch import Elasticsearch, helpers import json consumer KafkaConsumer( cdc-events.sales.orders, bootstrap_servers[kafka1:9092, kafka2:9092], group_ides-sync-worker, auto_offset_resetearliest ) es Elasticsearch([http://es1:9200]) def process_event(raw_msg): event json.loads(raw_msg.value) op event.get(op) # ccreate, uupdate, ddelete after event.get(after, {}) order_id after.get(order_id) or event.get(source, {}).get(key) if op d: es.delete(indexorders, idorder_id, ignore[404]) else: doc { order_id: after.get(order_id), user_id: after.get(user_id), goods_name: after.get(goods_name), amount: after.get(amount) } es.index(indexorders, idorder_id, documentdoc) for msg in consumer: process_event(msg)这里有个很重要的设计点ES 的文档_id直接用 Cassandra 的主键。好处是同一主键的数据无论重复消费多少次ES 里都只会有一份文档天然幂等不需要额外做去重这对消息管道来说省了非常多的事。4.3 ES 索引 mapping 的合理设计数据同步过去之前先把索引 mapping 设计好否则默认 mapping 会把所有字符串都识别成 text 类型查询时会出现精确值查不到的问题。我的订单索引 mapping 大概是这个样子的{ mappings: { properties: { order_id: { type: keyword }, user_id: { type: keyword }, goods_name: { type: text, analyzer: ik_max_word }, amount: { type: double }, status: { type: keyword }, created_at: { type: date } } } }几个设计要点order_id和user_id是精确匹配场景用keyword而不是text这样既能等值查询也能做聚合统计。goods_name要支持中文分词检索生产环境我一般装 IK 分词器analyzer 指定ik_max_word这样华为手机能拆成华为手机华为手机搜索结果更合理。amount是数值类型后续要对金额做 sum 聚合所以用double。created_at统一用date类型。Cassandra 里的时间戳和 ES 里的时间格式要做好转换建议统一存储 UTC 时间的 long 值或 ISO8601 字符串避免时区换算导致查询边界错位。5. 数据一致性这个环节决定成败5.1 管道里最容易出现的不一致情况不管双写还是 CDC最终一致是常态强一致反而不现实。我在项目里总结了几类高频不一致场景双写场景下Cassandra 成功、ES 超时或报错导致新增数据查不到。CDC 场景下消费端写入 ES 失败后又不断重试而 Cassandra 那边数据已经被更新二次ES 里还停在上一个版本。删除场景最容易被忽略Cassandra 里删除一行数据后CDC 会发出 delete 事件但如果消费端没有处理 delete 操作ES 里就永远残留这条已删除文档。Cassandra 的 tombstone 机制导致物理删除和逻辑删除并存如果 CDC 配置不当可能捕获不到删除操作。5.2 对账任务定时找出所有差异我在系统里设计了一个对账任务核心思想是定期比较 Cassandra 和 ES 两个数据源找出不一致的记录然后触发补偿。对账任务不是去全表比对那样效率太低。我采用增量时间窗口 抽查的方式假设当前时间是T对账任务扫描最近N分钟内 Cassandra 中更新的数据通过last_updated字段然后逐个去 ES 里查对应文档比较两个关键字段last_updated和version。比较逻辑可以做成这样def check_consistency(cass_row, es_doc): if es_doc is None: return MISSING_IN_ES if cass_row[version] ! es_doc[_source].get(version): return VERSION_MISMATCH if cass_row[last_updated] ! es_doc[_source].get(last_updated): return TIME_MISMATCH return OKversion字段是我在 Cassandra 表里额外维护的一个序号每次更新时自增。同步到 ES 后文档里也会带上这个 version。这样即使内容碰巧相同版本号变了也说明有一次更新没有正常同步。5.3 补偿执行策略对账发现问题后不能直接改 ES 就完事要走补偿流程。我的做法是把差异记录写入一张consistency_diff表由补偿任务逐条处理如果 Es 里缺文档就从 Cassandra 读最新数据重新写入如果版本不一致也以 Cassandra 为准重新覆盖如果 Cassandra 里已经是删除状态但 ES 里还有就执行一次 delete 操作。补偿任务要注意幂等和限流。幂等性靠的是以 Cassandra 为准、文档 ID 恒定限流是因为补偿任务往往集中爆发比如网络抖动恢复后如果一次性把所有差异都灌给 ES可能又把 ES 打挂所以我会给补偿任务加一个队列控制消费速率。5.4 顺序问题同一主键的更新乱序怎么办还有一个容易踩的深水区同一条订单数据被连续更新多次消息在 Kafka 里如果分到了不同分区消费端处理顺序就无法保证。比如 update(version2) 先被消费update(version3) 后到但内容可能覆盖旧数据。解决方案是生产端在发 Kafka 消息时按主键做 hash 分区保证同一个主键的所有变更事件进同一个分区。Cassandra 主键对应 Kafka 消息 key这样 partition 内天然有序。消费端在写入 ES 时再叠加一层乐观锁判断如果 ES 文档里的 version 大于等于当前事件里的 version直接丢弃否则才写入。两层保险基本能保证最终一致性。if es_doc and es_doc[_source].get(version, 0) event_version: # 丢弃过期事件 return es.index(indexorders, idorder_id, documentdoc)6. 查询优化与容量评估别等线上卡了再调6.1 ES 索引的生命周期管理做过一段时间 ES 的人都会明白一个索引从头用到尾是种奢望。订单、日志这类时序特征明显的数据我建议按天或按月建索引比如orders-2024-11-01、orders-2024-11-02。好处是数据过期可以直接删整个索引不用逐条 delete查询历史数据时可以按时间范围限定索引列表减少扫描分片数。ES 自带的 ILMIndex Lifecycle Management可以自动完成 rollover、压缩和删除。给索引挂一个别名设置策略{ policy: { phases: { hot: { actions: { rollover: { max_size: 50GB, max_age: 1d } } }, warm: { actions: { forcemerge: { max_num_segments: 1 } } }, delete: { actions: { delete: {} } } } } }这样索引超过 50GB 或 1 天就自动滚动新索引超过 30 天进入 warm 阶段做段合并超过 90 天自动删除。写入和查询都走别名应用层无感知。6.2 Cassandra 表结构设计如何反推同步链路同步链路跑得稳不稳和 Cassandra 表结构设计有很大关系。如果你把所有数据都塞进一个巨型表分区键设计不合理全量同步时扫描会很痛苦增量同步时 CDC 也会产生大量事件。我一般会在 Cassandra 侧按查询维度拆表典型的就是主表 查询索引表模式。主表按订单 ID 作为分区键负责接收写入另建一张按用户 ID 分区的表专门支撑查某用户所有订单的场景ES 里则承载所有非主键的检索需求。CDC 同步时两张表的事件都发到同一个 Kafka topic 或者不同 topic 分别消费最终在 ES 索引里归一化成同一个文档模型。这里要说一个很多人容易忽略的点Cassandra 里的UPDATE在底层其实也是插入一条新的 cell 数据overwrite并生成 tombstone 标记旧值。CDC 对这类的捕获逻辑和普通插入一致不需要特殊处理但因为 tombstone 会占用磁盘空间如果你的表频繁更新建议设置合理的gc_grace_seconds默认 10 天并定期执行nodetool compact清理。6.3 查询性能实测数据光说理论容易飘我贴一组来自我自己压测环境的数据6 节点 Cassandra 3 节点 ES单表 50 亿条订单数据让大家直观感受一下为什么查询必须走 ES查询场景Cassandra 方式响应耗时ES 方式响应耗时按订单 ID 精确查主键查询5mskeyword 等值查询3ms按用户 ID 查最近 100 条订单分区键 聚类列20ms范围查询 排序15ms商品名模糊搜索华为ALLOW FILTERING 全表扫12s还容易超时match 查询50ms按商品名分组统计金额不支持—terms 聚合 sum200ms按订单状态 时间范围过滤需要建额外索引表复杂bool 组合查询80ms模糊查询和聚合分析这两列就是决定性的差异。数据量过了亿级之后任何ALLOW FILTERING的查询都等于自杀式操作而 ES 走倒排索引再大的数据量也只是多几个分片的并行查询。这套组合里的查询响应速度从用户体验来说完全不在一个数量级。6.4 全量导入时 ES 的写入调优最后聊聊全量同步或大数据量灌入时ES 参数怎么调。默认配置下 ES 每秒钟刷盘一次还要实时维护副本如果同步峰值一小时写几百万条默认配置会非常吃力。我一般按以下方式调整在 bulk 写入期间把目标索引的副本数临时设为 0写入完成后改回 1。把refresh_interval从默认的1s调到30s甚至-1减少重建 segment 的开销。关闭translog.durability为async并调整sync_interval降低每次写入的磁盘等待。根据 record 大小估算 bulk 大小一般单次 bulk 控制在 5MB-15MB 之间用 8-16 个并发线程写入。写完以后执行一次POST /orders-*/_forcemerge?max_num_segments1把大量小 segment 合并成大 segment查询速度会有明显提升。当然这些调优参数是导入期的手段导入完成后一定要记得把副本数、refresh_interval 恢复成常规配置否则后续查询时索引会处于一种不健康的状态。最后再分享一点我的个人体会Cassandra 和 Elasticsearch 的整合技术难点从来不在怎么把数据写进 ES而在于你愿不愿意在设计阶段就把查询场景想清楚、把对账机制做扎实。项目上线第一周双写也好、CDC 也好看着都挺正常等跑了一个月数据量翻倍谁的对账任务在稳定兜底、谁的补偿逻辑在悄悄捞数据差距一下就拉开了。这套方案我已经在多个项目里反复验证过只要前期架构不歪后面维护起来是越来越省心的。