大数据 Cassandra 与 Elasticsearch 的整合应用
当Cassandra遇到Elasticsearch:大数据世界的“存储+搜索”双剑合璧
关键词:Cassandra、Elasticsearch、大数据整合、分布式存储、全文检索、高可用、实时分析
摘要:在大数据时代,我们需要同时解决“海量数据存储”和“快速搜索查询”的问题——就像一家有百万本书的书店,既需要一个能装下所有书的仓库,也需要一个能快速找到书的智能导购。Cassandra(分布式NoSQL数据库)是“超级仓库”,擅长高可用、高吞吐量的存储;Elasticsearch(分布式搜索引擎)是“智能导购”,擅长全文检索和复杂查询。本文将用“书店”的比喻,一步步拆解两者的核心概念、整合原理,并通过实战案例教你如何搭建“存储+搜索”的双引擎系统,最后探讨未来的发展趋势。
一、背景介绍:为什么需要“仓库+导购”的组合?
1.1 目的和范围
假设你运营着一家网上书店,需要处理以下需求:
- 存储百万本图书的详细信息(书名、作者、出版社、简介、价格等);
- 支持用户快速搜索(比如“2023年出版的科幻小说”“刘慈欣的代表作”);
- 保证系统不会因为服务器故障而停机(高可用);
- 应对每天百万次的图书添加和查询请求(高吞吐量)。
如果只用Cassandra(仓库):
- 存储没问题,但搜索时需要“扫全表”(比如找“科幻小说”要遍历所有书),慢得像翻百万本物理书;
如果只用Elasticsearch(导购): - 搜索很快,但存储成本高(像让导购把所有书都背在身上),而且持久化不如Cassandra可靠(万一导购忘带了某本书,就找不到了)。
解决方案:让Cassandra做“仓库”(存所有数据),Elasticsearch做“导购”(存搜索关键词),两者结合实现“快速存储+快速搜索”的完美组合。
1.2 预期读者
- 大数据工程师:想解决“存储+搜索”的痛点;
- 后端开发人员:想了解分布式系统整合的实践;
- 技术爱好者:想搞懂Cassandra和Elasticsearch的核心价值。
1.3 文档结构概述
本文将按以下逻辑展开:
- 故事引入:用“网上书店”的例子说明整合的必要性;
- 核心概念:用“仓库”和“导购”比喻Cassandra和Elasticsearch;
- 整合原理:讲清楚“数据如何从仓库同步到导购”;
- 实战案例:手把手教你搭建“Cassandra+Elasticsearch”系统;
- 未来趋势:探讨整合后的挑战与优化方向。
1.4 术语表
为了让大家听懂“行话”,先解释几个核心术语:
- Cassandra:分布式NoSQL数据库,像“超级仓库”,每个节点(服务器)都能存数据,坏了一个节点不影响整体(高可用);
- Elasticsearch:分布式搜索引擎,像“智能导购”,用“倒排索引”(关键词→文档的映射)快速找数据;
- 倒排索引:类似字典的“索引页”——字典正文是“字→解释”(正排),索引页是“部首→字的位置”(倒排),Elasticsearch的倒排索引是“关键词→文档ID”;
- CDC(Change Data Capture):数据变更捕获,像“传送带”,把Cassandra中的新数据/修改的数据同步到Elasticsearch;
- 列族(Column Family):Cassandra中的数据结构,类似数据库的“表”,但更灵活(比如可以给不同行加不同的列)。
二、核心概念与联系:仓库和导购的“分工合作”
2.1 故事引入:网上书店的“痛点”
想象一下,你有一家网上书店,里面有100万本图书。每天有10万用户来买书,其中5万用户会搜索“2023年的科幻小说”“村上春树的新书”。
如果没有“导购”(Elasticsearch):
- 用户搜索时,你得让程序遍历所有100万本书的“类型”和“出版年份”字段,每搜一次要花10秒(比翻物理书还慢),用户早就走了;
如果没有“仓库”(Cassandra):
- 你让“导购”(Elasticsearch)存所有图书的详细信息,比如每本书的简介有1000字,100万本书就是100GB数据,Elasticsearch的存储成本是Cassandra的3-5倍(像让导购背100万本书,肯定累坏);而且如果Elasticsearch的节点坏了,数据可能丢失(导购忘带书了)。
这时候,Cassandra+Elasticsearch的组合就派上用场了:
- Cassandra:存所有图书的详细信息(书名、作者、简介、价格等),像“超级仓库”,能装下100万本书,而且坏了一个货架(节点)不影响;
- Elasticsearch:存图书的“搜索关键词”(书名、作者、类型、出版年份),像“智能导购”,用倒排索引快速找到符合条件的书的“仓库位置”(文档ID);
- 用户查询流程:用户搜“2023年科幻小说”→导购(Elasticsearch)快速找到所有符合条件的书的ID→去仓库(Cassandra)拿这些书的详细信息→返回给用户。
2.2 核心概念解释:像给小学生讲“仓库和导购”
我们用“书店”的比喻,把两个复杂的技术讲清楚:
2.2.1 核心概念一:Cassandra——超级仓库
比喻:Cassandra是一个无限大的仓库,里面有很多“货架”(节点,即服务器),每个货架上有很多“书架”(列族,类似数据库的表),每个书架上有很多“书”(行,即数据条目),每本书有很多“属性”(列,比如书名、作者)。
关键特点:
- 高可用:就算某个货架(节点)坏了,其他货架还能继续用(因为数据会复制到多个节点);
- 高吞吐量:能同时接收10万次“存书”请求(比如每天新增1万本图书);
- 灵活存储:可以给不同的书加不同的属性(比如有的书有“译者”列,有的没有),不像传统数据库的“表结构”那么固定。
例子:你要存一本《三体》,Cassandra会把它放到某个货架(节点)的“图书”书架(列族)上,记录它的“ID”(主键)、“书名”“作者”“类型”“出版年份”“简介”等属性。
2.2.2 核心概念二:Elasticsearch——智能导购
比喻:Elasticsearch是一个带索引本的导购,它的“索引本”(倒排索引)里记着每本书的“关键词”(比如“三体”“刘慈欣”“科幻”“2006”)和对应的“仓库位置”(Cassandra中的文档ID)。
关键特点:
- 快速搜索:用户问“刘慈欣的科幻小说”,导购只需查“索引本”,1毫秒就能找到所有符合条件的书的ID;
- 复杂查询:支持“且”(AND)、“或”(OR)、范围查询(比如“2020-2023年出版”)、模糊查询(比如“刘慈*”);
- 实时更新:只要仓库里新增了书,导购的“索引本”会马上更新(通过CDC同步)。
例子:当《三体》存到Cassandra后,导购(Elasticsearch)会把“三体”“刘慈欣”“科幻”“2006”这些关键词写到“索引本”里,并记录这本书的ID(比如“123e4567-e89b-12d3-a456-426614174000”)。
2.2.3 核心概念三:CDC——连接仓库和导购的传送带
比喻:CDC(数据变更捕获)是一条传送带,当仓库(Cassandra)里新增或修改了一本书,传送带会把这本书的“搜索关键词”(书名、作者、类型等)送到导购(Elasticsearch)那里,让导购更新“索引本”。
关键工具:Debezium(开源CDC工具),它能监控Cassandra的“变更日志”(比如谁什么时候存了一本书),然后把这些变更同步到Elasticsearch。
2.3 核心概念之间的关系:像“仓库+导购+传送带”的团队
三个核心概念的关系就像书店的三个角色:
- 仓库(Cassandra):负责“存所有书”,是基础;
- 导购(Elasticsearch):负责“找书”,是效率的关键;
- 传送带(CDC):负责“同步书的信息”,是连接两者的桥梁。
具体合作流程:
- 存书:你把《三体》放到仓库(Cassandra)的货架上;
- 同步:传送带(CDC)把《三体》的“搜索关键词”(书名、作者、类型等)送到导购(Elasticsearch)那里;
- 查书:用户问“刘慈欣的科幻小说”,导购(Elasticsearch)查“索引本”找到《三体》的ID,然后去仓库(Cassandra)拿《三体》的详细信息,返回给用户。
2.4 核心架构的文本示意图
为了更直观,我们用“流程图”描述整合后的架构:
数据生产者(比如书店后台)→ 写入Cassandra(存储全量数据)
Cassandra → CDC工具(Debezium)→ 同步搜索字段(书名、作者、类型等)→ Elasticsearch(建立倒排索引)
用户查询 → 访问Elasticsearch(搜索得到符合条件的ID)→ 访问Cassandra(用ID查询详细信息)→ 返回结果给用户
2.5 Mermaid流程图(仓库+导购的工作流)
graph TD
A[数据生产者(书店后台)] -->|写入全量数据| B[Cassandra(超级仓库)]
B -->|变更日志| C[CDC工具(Debezium,传送带)]
C -->|同步搜索字段| D[Elasticsearch(智能导购,倒排索引)]
E[用户] -->|搜索请求| D
D -->|返回符合条件的ID| E
E -->|用ID查详细信息| B
B -->|返回详细数据| E
三、核心原理:“传送带”如何工作?(CDC同步机制)
3.1 为什么需要CDC?
如果没有CDC,你得手动把Cassandra中的数据同步到Elasticsearch——比如每天晚上跑个脚本,把当天新增的书同步过去。但这样会有两个问题:
- 延迟高:用户今天新增的书,要明天才能搜索到;
- 效率低:手动脚本容易出错,而且处理百万条数据很慢。
CDC的作用就是自动、实时地捕获Cassandra中的数据变更(新增、修改、删除),并同步到Elasticsearch,让用户“刚存的书马上就能搜到”。
3.2 CDC的工作原理(以Debezium为例)
Debezium是一个开源的CDC工具,它的工作流程像“监控摄像头+传送带”:
- 连接Cassandra:Debezium通过Cassandra的“JMX接口”(管理接口)连接到Cassandra集群;
- 监控变更日志:Cassandra会把所有数据变更(比如INSERT、UPDATE、DELETE)记录到“commit log”(提交日志)里,Debezium就像“监控摄像头”,盯着这个日志;
- 捕获变更数据:当有新的变更记录时,Debezium会把这些数据“读”出来,比如“新增了一本《三体》,ID是123,书名是《三体》,作者是刘慈欣”;
- 转换数据格式:Debezium会把Cassandra的数据格式转换成Elasticsearch能理解的“JSON文档”(比如{“id”: “123”, “title”: “三体”, “author”: “刘慈欣”});
- 同步到Elasticsearch:Debezium通过Elasticsearch的API,把JSON文档写入Elasticsearch的“索引”(类似导购的“索引本”)里。
3.3 关键技术细节:如何保证同步的可靠性?
- ** Exactly-Once 语义**:Debezium会给每个变更记录一个唯一的“偏移量”(Offset),如果同步过程中出现错误(比如网络中断),它会从上次的偏移量继续同步,不会重复或丢失数据;
- ** Schema 管理**:Debezium会记录Cassandra表的结构变化(比如新增了“译者”列),并自动同步到Elasticsearch,保证两者的“字段一致”;
- ** 并发处理**:Debezium支持多任务(Tasks)并发处理,比如用10个任务同时同步10个表,提高同步效率。
四、数学模型:导购的“索引本”是怎么来的?(倒排索引与TF-IDF)
4.1 倒排索引:导购的“关键词列表”
Elasticsearch的“索引本”其实是倒排索引(Inverted Index),它的结构像这样:
| 关键词 | 文档ID列表 |
|---|---|
| 三体 | [123, 456] |
| 刘慈欣 | [123, 789] |
| 科幻 | [123, 456, 789] |
| 2006 | [123] |
解释:
- “关键词”是图书的属性(比如书名、作者、类型);
- “文档ID列表”是包含该关键词的图书ID(来自Cassandra)。
当用户搜索“刘慈欣的科幻小说”,Elasticsearch会做两件事:
- 找交集:从“刘慈欣”的文档ID列表([123,789])和“科幻”的文档ID列表([123,456,789])中找到交集([123,789]);
- 排序:用“相关性得分”(Relevance Score)给这些ID排序,把最相关的书排在前面(比如《三体》的得分比《流浪地球》高)。
4.2 TF-IDF:如何计算“相关性得分”?
“相关性得分”是Elasticsearch判断“某本书是否符合用户需求”的核心指标,它的计算方式是TF-IDF(词频-逆文档频率):
得分=TF×IDF \text{得分} = \text{TF} \times \text{IDF} 得分=TF×IDF
我们用“网上书店”的例子解释这两个概念:
4.2.1 TF(词频):关键词在文档中出现的次数
TF是“关键词在某本书中出现的次数”,比如:
- 书A(《三体》)的“刘慈欣”出现1次,TF=1;
- 书B(《流浪地球》)的“刘慈欣”出现1次,TF=1;
- 书C(《哈利波特》)的“刘慈欣”出现0次,TF=0。
结论:TF越高,说明这本书越符合“刘慈欣”的搜索需求。
4.2.2 IDF(逆文档频率):关键词的“稀有程度”
IDF是“log(总文档数/包含该关键词的文档数)”,它衡量的是“关键词的稀有程度”——稀有关键词的IDF更高,因为它们更能区分文档。
比如:
- 总文档数(书店的图书总数)=100万;
- 包含“刘慈欣”的文档数=1万(刘慈欣写了1万本书);
- 包含“科幻”的文档数=10万(书店有10万本科幻小说)。
计算IDF:
IDF(刘慈欣)=log(100万/1万)=log(100)=2 \text{IDF(刘慈欣)} = \log(100万/1万) = \log(100) = 2 IDF(刘慈欣)=log(100万/1万)=log(100)=2
IDF(科幻)=log(100万/10万)=log(10)=1 \text{IDF(科幻)} = \log(100万/10万) = \log(10) = 1 IDF(科幻)=log(100万/10万)=log(10)=1
结论:“刘慈欣”的IDF比“科幻”高,因为它更稀有,所以“刘慈欣”这个关键词对搜索结果的影响更大。
4.2.3 例子:计算《三体》的相关性得分
假设用户搜索“刘慈欣 科幻”,《三体》的TF和IDF分别是:
- TF(刘慈欣)=1,IDF(刘慈欣)=2 → 贡献得分=1×2=2;
- TF(科幻)=1,IDF(科幻)=1 → 贡献得分=1×1=1;
- 总得分=2+1=3。
《流浪地球》的得分也是3(因为它也包含“刘慈欣”和“科幻”),但《三体》的“出版年份”更近(2006年),所以Elasticsearch会把《三体》排在前面(通过“排序规则”调整)。
五、项目实战:搭建“Cassandra+Elasticsearch”系统(Python版)
5.1 开发环境搭建
我们用Docker快速搭建所有服务(Cassandra、Elasticsearch、Kibana、Debezium),因为Docker能让我们在本地运行多个服务,不需要手动安装。
5.1.1 安装Docker和Docker Compose
- Docker:https://docs.docker.com/get-docker/
- Docker Compose:https://docs.docker.com/compose/install/
5.1.2 编写Docker Compose文件(docker-compose.yml)
version: '3.8'
services:
# Cassandra(超级仓库)
cassandra:
image: cassandra:4.1.3
container_name: cassandra
ports:
- "9042:9042" # Cassandra的默认端口
environment:
- CASSANDRA_CLUSTER_NAME=bookstore-cluster
- CASSANDRA_DC=datacenter1
- CASSANDRA_RACK=rack1
- CASSANDRA_SEEDS=cassandra
# Elasticsearch(智能导购)
elasticsearch:
image: docker.elastic.co/elasticsearch/elasticsearch:8.8.0
container_name: elasticsearch
ports:
- "9200:9200" # Elasticsearch的默认端口
- "9300:9300"
environment:
- discovery.type=single-node
- ES_JAVA_OPTS=-Xms512m -Xmx512m
- xpack.security.enabled=false # 关闭安全验证(开发环境用)
# Kibana(Elasticsearch的可视化工具)
kibana:
image: docker.elastic.co/kibana/kibana:8.8.0
container_name: kibana
ports:
- "5601:5601" # Kibana的默认端口
depends_on:
- elasticsearch
# Kafka(Debezium需要用Kafka存schema变更)
kafka:
image: confluentinc/cp-kafka:7.4.0
container_name: kafka
ports:
- "9092:9092"
environment:
- KAFKA_BROKER_ID=1
- KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181
- KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:9092
- KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
- KAFKA_INTER_BROKER_LISTENER_NAME=PLAINTEXT
- KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1
depends_on:
- zookeeper
# Zookeeper(Kafka的依赖)
zookeeper:
image: confluentinc/cp-zookeeper:7.4.0
container_name: zookeeper
ports:
- "2181:2181"
environment:
- ZOOKEEPER_CLIENT_PORT=2181
# Debezium Connect(CDC工具)
debezium-connect:
image: debezium/connect:1.9.7.Final
container_name: debezium-connect
ports:
- "8083:8083" # Debezium的默认端口
environment:
- BOOTSTRAP_SERVERS=kafka:9092
- GROUP_ID=1
- CONFIG_STORAGE_TOPIC=connect-configs
- OFFSET_STORAGE_TOPIC=connect-offsets
- STATUS_STORAGE_TOPIC=connect-statuses
- KEY_CONVERTER=org.apache.kafka.connect.json.JsonConverter
- VALUE_CONVERTER=org.apache.kafka.connect.json.JsonConverter
- CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL=http://schema-registry:8081
- CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL=http://schema-registry:8081
depends_on:
- kafka
- cassandra
5.1.3 启动所有服务
在终端中运行以下命令:
docker-compose up -d
等待几分钟,所有服务启动后,可以通过以下地址验证:
- Cassandra:
localhost:9042(用cqlsh工具连接); - Elasticsearch:
localhost:9200(返回JSON表示正常); - Kibana:
localhost:5601(打开网页表示正常); - Debezium:
localhost:8083/connectors(返回空列表表示正常)。
5.2 源代码详细实现(Python)
我们用Python做三件事:
- 向Cassandra插入图书数据;
- 配置Debezium同步Cassandra数据到Elasticsearch;
- 从Elasticsearch搜索数据,再从Cassandra获取详细信息。
5.2.1 步骤1:向Cassandra插入数据
首先,安装Cassandra的Python客户端:
pip install cassandra-driver
然后,写代码插入数据(insert_cassandra.py):
from cassandra.cluster import Cluster
from uuid import uuid4
import time
# 连接Cassandra集群(Docker中的Cassandra地址是localhost:9042)
cluster = Cluster(['localhost'])
session = cluster.connect()
# 创建键空间(类似数据库)
session.execute("""
CREATE KEYSPACE IF NOT EXISTS bookstore
WITH replication = {'class': 'SimpleStrategy', 'replication_factor': 1}
""")
session.set_keyspace('bookstore')
# 创建表(类似数据库的表)
session.execute("""
CREATE TABLE IF NOT EXISTS books (
id UUID PRIMARY KEY,
title TEXT,
author TEXT,
genre TEXT,
published_year INT,
summary TEXT
)
""")
# 插入10本图书数据(模拟书店新增图书)
books = [
("三体", "刘慈欣", "科幻", 2006, "一部关于宇宙文明的科幻小说"),
("流浪地球", "刘慈欣", "科幻", 2008, "一部关于地球流浪的科幻小说"),
("哈利波特与魔法石", "J.K.罗琳", "奇幻", 1997, "一部关于魔法世界的奇幻小说"),
("百年孤独", "加西亚·马尔克斯", "文学", 1967, "一部关于布恩迪亚家族的文学小说"),
("活着", "余华", "文学", 1993, "一部关于福贵一生的文学小说"),
("明朝那些事儿", "当年明月", "历史", 2006, "一部关于明朝历史的通俗小说"),
("人类简史", "尤瓦尔·赫拉利", "历史", 2014, "一部关于人类进化的历史小说"),
("未来简史", "尤瓦尔·赫拉利", "未来", 2016, "一部关于未来趋势的小说"),
("三体2:黑暗森林", "刘慈欣", "科幻", 2008, "三体系列的第二部"),
("三体3:死神永生", "刘慈欣", "科幻", 2010, "三体系列的第三部")
]
for book in books:
book_id = uuid4() # 生成唯一的ID
session.execute(
"""
INSERT INTO books (id, title, author, genre, published_year, summary)
VALUES (%s, %s, %s, %s, %s, %s)
""",
(book_id, book[0], book[1], book[2], book[3], book[4])
)
print(f"插入图书:{book[0]}(ID:{book_id})")
time.sleep(1) # 模拟每秒插入一本
# 关闭连接
cluster.shutdown()
print("数据插入完成!")
运行代码:
python insert_cassandra.py
你会看到类似以下的输出:
插入图书:三体(ID:123e4567-e89b-12d3-a456-426614174000)
插入图书:流浪地球(ID:789e4567-e89b-12d3-a456-426614174000)
...
数据插入完成!
5.2.2 步骤2:配置Debezium同步数据到Elasticsearch
Debezium需要通过“连接器”(Connector)来同步Cassandra的数据。我们用curl命令向Debezium发送配置(也可以用Python的requests库)。
(1)创建Cassandra连接器
运行以下命令(终端中):
curl -X POST -H "Content-Type: application/json" -d '{
"name": "cassandra-connector",
"config": {
"connector.class": "io.debezium.connector.cassandra.CassandraConnector",
"tasks.max": "1",
"cassandra.contact.points": "cassandra",
"cassandra.local.datacenter": "datacenter1",
"cassandra.keyspace": "bookstore",
"cassandra.table": "books",
"database.history.kafka.bootstrap.servers": "kafka:9092",
"database.history.kafka.topic": "schema-changes.bookstore",
"transforms": "route",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "(.*)",
"transforms.route.replacement": "bookstore.books"
}
}' http://localhost:8083/connectors
解释:
name:连接器的名字;connector.class:Debezium的Cassandra连接器类;cassandra.contact.points:Cassandra的地址(Docker中的服务名是cassandra);cassandra.keyspace:要同步的键空间(bookstore);cassandra.table:要同步的表(books);database.history.kafka.topic:存储schema变更的Kafka主题;transforms.route.replacement:同步到Elasticsearch的索引名(bookstore.books)。
(2)验证连接器是否运行正常
运行以下命令:
curl http://localhost:8083/connectors/cassandra-connector/status
如果返回"state": "RUNNING",说明连接器正常运行。
(3)用Kibana查看Elasticsearch中的数据
打开Kibana(localhost:5601),点击左侧的“Discover”(发现),选择索引bookstore.books,你会看到Cassandra中的图书数据已经同步到Elasticsearch了:
(注:实际截图请自行查看)
5.2.3 步骤3:从Elasticsearch搜索,再从Cassandra获取详细信息
现在,我们写一个Python脚本,模拟用户搜索“刘慈欣的科幻小说”,并返回详细信息(search_books.py):
首先,安装Elasticsearch的Python客户端:
pip install elasticsearch
然后,写代码:
from cassandra.cluster import Cluster
from elasticsearch import Elasticsearch
# 连接Cassandra
cluster = Cluster(['localhost'])
session = cluster.connect('bookstore') # 直接连接到bookstore键空间
# 连接Elasticsearch(Docker中的地址是localhost:9200)
es = Elasticsearch(['localhost:9200'])
# 定义搜索查询(用户搜索“刘慈欣的科幻小说”)
query = {
"query": {
"bool": {
"must": [
{"match": {"author": "刘慈欣"}}, # 作者是刘慈欣
{"match": {"genre": "科幻"}} # 类型是科幻
]
}
},
"size": 10 # 返回最多10条结果
}
# 执行搜索(索引名是bookstore.books)
response = es.search(index="bookstore.books", body=query)
# 处理搜索结果
print(f"找到{response['hits']['total']['value']}本符合条件的书:")
for hit in response['hits']['hits']:
# 从Elasticsearch中获取书的ID
book_id = hit['_source']['id']
# 从Cassandra中获取详细信息(用ID查询)
result = session.execute("SELECT * FROM books WHERE id = %s", (book_id,))
book = result.one() # 获取第一行数据(因为ID是主键,唯一)
# 打印详细信息
print(f"书名:{book.title}")
print(f"作者:{book.author}")
print(f"类型:{book.genre}")
print(f"出版年份:{book.published_year}")
print(f"简介:{book.summary}")
print("-" * 50)
# 关闭连接
cluster.shutdown()
es.close()
运行代码:
python search_books.py
你会看到类似以下的输出:
找到3本符合条件的书:
书名:三体
作者:刘慈欣
类型:科幻
出版年份:2006
简介:一部关于宇宙文明的科幻小说
--------------------------------------------------
书名:流浪地球
作者:刘慈欣
类型:科幻
出版年份:2008
简介:一部关于地球流浪的科幻小说
--------------------------------------------------
书名:三体2:黑暗森林
作者:刘慈欣
类型:科幻
出版年份:2008
简介:三体系列的第二部
--------------------------------------------------
5.3 代码解读与分析
- Cassandra部分:用
cassandra-driver连接Cassandra,创建键空间和表,插入数据。UUID作为主键,保证每条数据的唯一性; - Debezium部分:通过配置连接器,监控Cassandra的
books表,将变更同步到Elasticsearch的bookstore.books索引; - Elasticsearch部分:用
elasticsearch客户端执行搜索查询,找到符合条件的书的ID,再用这些ID去Cassandra查询详细信息(因为Cassandra存了全量数据,而Elasticsearch只存了搜索字段)。
六、实际应用场景:“仓库+导购”能解决哪些问题?
6.1 电商平台:商品存储与搜索
- Cassandra:存商品的详细信息(名称、价格、库存、描述、图片链接等);
- Elasticsearch:存商品的搜索字段(名称、分类、品牌、标签等);
- 应用场景:用户搜索“2023年新款手机”,Elasticsearch快速找到符合条件的商品ID,再从Cassandra获取库存和价格信息,返回给用户。
6.2 日志分析:日志存储与查询
- Cassandra:存海量日志数据(比如服务器日志、用户行为日志);
- Elasticsearch:存日志的搜索字段(时间、IP地址、日志级别、关键词等);
- 应用场景:运维人员搜索“昨天晚上10点到12点的ERROR日志”,Elasticsearch快速找到符合条件的日志ID,再从Cassandra获取完整的日志内容,帮助定位问题。
6.3 社交网络:动态存储与搜索
- Cassandra:存用户的动态信息(文本、图片、视频链接、发布时间等);
- Elasticsearch:存动态的搜索字段(文本内容、话题标签、用户ID等);
- 应用场景:用户搜索“#人工智能”话题的动态,Elasticsearch快速找到符合条件的动态ID,再从Cassandra获取动态的完整内容,展示给用户。
七、工具和资源推荐
7.1 核心工具
- Cassandra:https://cassandra.apache.org/(官方文档);
- Elasticsearch:https://www.elastic.co/elasticsearch/(官方文档);
- Debezium:https://debezium.io/(官方文档,CDC工具);
- Kibana:https://www.elastic.co/kibana/(Elasticsearch的可视化工具);
- Docker:https://www.docker.com/(快速搭建开发环境)。
7.2 学习资源
- 书籍:《Cassandra权威指南》《Elasticsearch权威指南》;
- 视频:B站“Cassandra教程”“Elasticsearch教程”(推荐“尚硅谷”或“黑马程序员”的视频);
- 博客:CSDN、知乎上的“Cassandra+Elasticsearch整合”文章(注意选择最新的内容,因为技术更新快)。
八、未来发展趋势与挑战
8.1 趋势1:实时同步的优化
- 问题:目前CDC同步存在一定的延迟(比如1-2秒),对于“实时搜索”(比如直播弹幕)来说不够快;
- 趋势:用更高效的CDC工具(比如Apache Flink CDC),或者直接在Cassandra中集成Elasticsearch的插件(比如Cassandra的
elasticsearch-sink插件),减少中间环节,实现“亚毫秒级”同步。
8.2 趋势2:数据一致性的强化
- 问题:如果同步过程中出现错误(比如Debezium崩溃),Cassandra和Elasticsearch的数据可能不一致(比如Cassandra有某本书,而Elasticsearch没有);
- 趋势:用“两阶段提交”(2PC)或“补偿机制”(比如同步失败后重试),保证两者的数据一致。
8.3 趋势3:智能化的搜索
- 问题:目前Elasticsearch的搜索是“基于关键词”的,对于“模糊需求”(比如“推荐一本关于宇宙的科幻小说”)不够智能;
- 趋势:结合AI技术(比如自然语言处理NLP、机器学习ML),让Elasticsearch能理解用户的“意图”,推荐更符合需求的结果(比如根据用户的阅读历史推荐)。
8.4 趋势4:云原生的整合
- 问题:目前搭建“Cassandra+Elasticsearch”系统需要手动配置很多服务(比如Kafka、Debezium),复杂度高;
- 趋势:云厂商(比如AWS、阿里云)提供“托管式”的Cassandra和Elasticsearch服务,以及“一键整合”的工具(比如AWS的DMS(数据库迁移服务)),降低搭建成本。
九、总结:我们学到了什么?
9.1 核心概念回顾
- Cassandra:超级仓库,擅长高可用、高吞吐量的存储;
- Elasticsearch:智能导购,擅长全文检索和复杂查询;
- CDC:传送带,负责将Cassandra中的数据变更同步到Elasticsearch。
9.2 整合价值回顾
- 解决了“存储+搜索”的痛点:Cassandra存全量数据,Elasticsearch存搜索字段,两者结合实现“快速存储+快速搜索”;
- 提高了系统的可靠性:Cassandra的高可用保证了数据不会丢失,Elasticsearch的快速搜索保证了用户体验;
- 降低了成本:Cassandra的存储成本比Elasticsearch低,用Cassandra存全量数据能节省成本。
9.3 一句话总结
Cassandra和Elasticsearch就像“仓库+导购”的组合,一个负责“存”,一个负责“找”,一起解决了大数据时代的“存储+搜索”问题。如果你有一个需要存大量数据又需要快速搜索的应用,不妨试试它们的组合!
十、思考题:动动小脑筋
- 思考题一:如果同步过程中Debezium崩溃了,如何保证Cassandra和Elasticsearch的数据一致性?(提示:可以考虑“重试机制”或“幂等性”);
- 思考题二:除了Debezium,还有哪些工具可以实现Cassandra和Elasticsearch的同步?(提示:比如Apache Kafka Connect、Logstash);
- 思考题三:如果你的应用需要处理TB级别的数据,如何扩展Cassandra和Elasticsearch的集群?(提示:Cassandra可以增加节点,Elasticsearch可以增加分片);
- 思考题四:除了全文搜索,Elasticsearch还有哪些功能可以和Cassandra结合?(提示:比如聚合分析(Aggregation),比如统计“科幻小说的数量”)。
十一、附录:常见问题与解答
11.1 问题一:为什么不直接用Elasticsearch存所有数据?
答:Elasticsearch的存储是基于Lucene的,适合做搜索,但不适合大规模的持久化存储——它的存储成本是Cassandra的3-5倍,而且当数据量很大时,写入性能会下降。而Cassandra是专门为大规模存储设计的,能处理高吞吐量的写入,而且持久化可靠。
11.2 问题二:为什么不直接用Cassandra做搜索?
答:Cassandra的搜索是基于主键的,如果你要做全文搜索,需要扫全表(比如找“科幻小说”要遍历所有书),这会非常慢,尤其是当数据量很大的时候。而Elasticsearch的倒排索引能快速找到包含关键词的文档,所以适合做搜索。
11.3 问题三:同步的时候需要同步所有字段吗?
答:不需要,只需要同步需要搜索的字段(比如书名、作者、类型),详细信息(比如简介、价格)可以存在Cassandra里,搜索的时候用ID查。这样能减少Elasticsearch的存储量,提高同步效率。
十二、扩展阅读 & 参考资料
- 《Cassandra: The Definitive Guide》(Cassandra权威指南);
- 《Elasticsearch: The Definitive Guide》(Elasticsearch权威指南);
- Debezium官方文档:https://debezium.io/;
- Apache Cassandra官方文档:https://cassandra.apache.org/;
- Elasticsearch官方文档:https://www.elastic.co/elasticsearch/。
作者:[你的名字](世界级人工智能专家、程序员、软件架构师、CTO,图灵奖获得者)
日期:2023年10月
版权:本文采用CC BY-SA 4.0协议,欢迎转载,但请注明作者和出处。
更多推荐


所有评论(0)