基于 Cassandra、SolrCloud、Zookeeper、Kafka 和 Python 微服务的搜索系统搭建实践
本方案假设 Linux 服务器环境(CentOS 或 Ubuntu)已安装适当版本的 Java JDK(Cassandra、Solr、Kafka 都依赖 Java)。多台服务器部署方案如下:
- Zookeeper 集群:3 台节点,用于协调 Kafka 和 SolrCloud(IP 范围 10.201.X.X,客户端端口 2181)。
- Kafka 集群:3 台或以上,使用 Zookeeper 集群协调。Kafka 负责接收和传递产品数据更新消息。
- Cassandra:1 台或多台,根据数据量和容错需求而定。存储产品数据(单节点 IP 10.201.X.X)。
- SolrCloud 集群:3 台 Solr 节点组成搜索索引集群(IP 分别为 10.201.X.X)。
- Python 微服务:
- 数据同步服务:消费 Kafka 消息,将新产品数据写入 Cassandra,更新 Solr 索引。
- 搜索 API 服务:提供 REST 接口供前端查询,通过 Solr 获取搜索结果(也可在服务层做权限控制和聚合)。
各组件通过内部网络通信,防火墙需放行必要端口(Cassandra 9042、Solr 8983、Zookeeper 2181、Kafka 9092 及微服务配置端口)。示例中使用 10.201.X.X 作为占位符;生产环境请替换为实际地址。
安装与配置 Cassandra
安装 Cassandra
从 Apache 官方下载 Cassandra(本例使用 3.11.12 版本),上传到目标服务器 10.201.X.X 并解压:
$ wget https://downloads.apache.org/cassandra/3.11.12/apache-cassandra-3.11.12-bin.tar.gz
$ tar -zxvf apache-cassandra-3.11.12-bin.tar.gz
$ cd apache-cassandra-3.11.12
Cassandra 无需额外编译。进入 Cassandra 目录可通过 bin/cassandra 启动节点。首次启动前需修改配置以确保节点通信正常。
配置 Cassandra 节点
编辑 conf/cassandra.yaml,按网络环境调整:
- cluster_name:集群名称,如 cluster_name: "SearchCluster"。
- listen_address:节点监听地址,设为本机内网 IP,如 listen_address: 10.201.X.X。
- rpc_address:控制 CQL 服务绑定地址。设为节点 IP,如 rpc_address: 10.201.X.X(Cassandra 3.x 中 rpc_address 控制 CQL)。
- seed_provider:种子节点列表,用于集群引导。单节点测试时设为本机 IP;多节点集群时列出初始种子节点 IP:
seed_provider:
- class_name: org.apache.cassandra.locator.SimpleSeedProvider
parameters:
- seeds: "10.201.X.X,10.201.X.X"
至少包含集群中一个节点 IP 作为种子。
- 其他参数:按需调整内存和垃圾回收设置。默认配置通常可用;生产环境根据可用内存调整 -Xmx 等参数。
启动 Cassandra:
$ bin/cassandra -R # 后台启动 Cassandra (-R 可去除超级用户模式限制)
初次启动可能需等待片刻完成初始化。使用 CQL Shell 验证连接:
$ bin/cqlsh 10.201.X.X 9042 # 连接到本机的 Cassandra CQL 服务
Connected to SearchCluster at 10.201.X.X:9042.
cqlsh>
连接成功后创建 Keyspace 和表。例如创建 search keyspace:
CREATE KEYSPACE search
WITH replication = {'class': 'SimpleStrategy', 'replication_factor': '1'}
AND durable_writes = true;
注意: 上述使用简单策略和副本因子 1,仅适用测试。生产集群应根据数据中心拓扑使用 NetworkTopologyStrategy 并设定合适副本数。
设计 Cassandra 表:按查询模式组织数据。Cassandra 适合存储宽行数据,模式设计应围绕查询需求。在 search keyspace 下创建多张表:
- 产品主表(goods):存储产品主要字段,主键为产品 ID,包含名称、描述、分类等。
- 属性表:按属性类型拆分,主键为产品 ID,各表存储产品特定属性集合,含版本号记录更新。
- 维度表(category 等):按业务需要创建,如按类目、品牌、促销等维度组织数据。
搜索主要由 Solr 执行查询,Cassandra 充当持久化存储和数据源;通常搜索时不直接访问 Cassandra,而是预先将数据导入 Solr。
安装与配置 Zookeeper
Zookeeper 在本方案中有两个角色:Kafka 的协调服务和 SolrCloud 的配置管理。我们搭建由 3 台节点组成的集群,保证高可用。
安装 Zookeeper
下载 Zookeeper(如 3.7.0 版本)并在每台服务器上解压:
$ wget https://downloads.apache.org/zookeeper/zookeeper-3.7.0/apache-zookeeper-3.7.0-bin.tar.gz
$ tar -zxvf apache-zookeeper-3.7.0-bin.tar.gz
$ mv apache-zookeeper-3.7.0-bin /opt/app/zookeeper
假设三台服务器 IP 分别为 10.201.X.X、10.201.X.X、10.201.X.X。
配置 Zookeeper
在每台节点的 conf 目录下复制模板:
$ cp conf/zoo_sample.cfg conf/zoo.cfg
编辑 zoo.cfg,关注以下配置:
- dataDir=/path/to/zookeeper/data:数据存储目录,如 /opt/app/zookeeper/data,确保目录存在且有读写权限。
- clientPort=2181:客户端连接端口,默认 2181,一般保持不变。
- 集群配置:在文件末尾添加:
server.1=10.201.X.X:2888:3888
server.2=10.201.X.X:2888:3888
server.3=10.201.X.X:2888:3888
server.N 的 N 对应每台服务器的唯一编号。2888 是集群通信端口,3888 为选举端口。
- 在每台节点的 dataDir 中创建 myid 文件,内容为该节点编号(节点 1 写入"1",节点 2 写入"2",等等)。
启动 Zookeeper 集群
在每台服务器上运行:
$ ./bin/zkServer.sh start
启动后用 zkServer.sh status 检查状态,应看到一台 leader、其他为 follower。验证 ZK 集群可用性,用 zkCli 连接任意节点:
$ ./bin/zkCli.sh -server 10.201.X.X:2181
测试创建节点:
[zk: 10.201.X.X:2181(CONNECTED) 0] create /solr_demo "test"
成功创建说明 ZK 正常工作。后续将在 ZK 的 /solr_demo 路径存储 SolrCloud 配置数据。
安装与配置 Kafka
安装 Kafka
从 Kafka 官网或 Apache 存档下载(本例使用 Kafka 2.x):
$ wget https://downloads.apache.org/kafka/2.7.0/kafka_2.12-2.7.0.tgz
$ tar -zxvf kafka_2.12-2.7.0.tgz
$ mv kafka_2.12-2.7.0 /opt/app/kafka
在每台 Kafka 服务器上解压。Kafka 需要 JDK 1.8+。
配置 Kafka Broker
编辑 config/server.properties:
- broker.id:集群中每个节点需唯一(整数)。第一台设为 0,第二台 1,第三台 2,等等。
- listeners 和 advertised.listeners:配置监听地址和对外通告地址。单机默认 PLAINTEXT://:9092。多网卡或容器环境需设置 advertised.listeners 为实际可达地址,如 PLAINTEXT://10.201.X.X:9092。
- zookeeper.connect:Zookeeper 集群地址列表:
zookeeper.connect=10.201.X.X:2181,10.201.X.X:2181,10.201.X.X:2181
SolrCloud 若使用 ZK 的 /solr_demo 子路径,Kafka 可直接连根目录;Kafka 会在 ZK 上创建 /brokers 等路径。
- log.dirs:Kafka 存储消息的目录,如 /opt/app/kafka/logs。
根据需要调整内存等参数,但默认配置适合初始测试。
启动 Kafka
依次在每台 Kafka 节点启动 broker:
$ ./bin/kafka-server-start.sh -daemon config/server.properties
-daemon 参数使进程后台运行。启动后 Kafka 在 Zookeeper 注册自己。可通过 Zookeeper 客户端查看 /brokers/ids 节点下是否有对应的 broker ID。
创建主题
根据业务需要创建 Kafka 主题传递产品数据更新。例如创建 product_update 主题,副本因子 2,分区数 3:
$ ./bin/kafka-topics.sh --create --topic product_update --partitions 3 --replication-factor 2 --bootstrap-server 10.201.X.X:9092
(--bootstrap-server 指定任意 Kafka 节点地址)。验证创建成功,列出当前主题:
$ ./bin/kafka-topics.sh --list --bootstrap-server 10.201.X.X:9092
Kafka 集群至此就绪。
安装与配置 SolrCloud
Solr 是搜索索引核心组件。使用 SolrCloud 模式部署以支持分片和高可用(本例使用 Solr 7.7.3)。
安装 Solr
从 Solr 官方下载对应版本(zip 或 tgz),在每台节点解压:
$ wget https://archive.apache.org/dist/lucene/solr/7.7.3/solr-7.7.3.tgz
$ tar -zxvf solr-7.7.3.tgz
$ mv solr-7.7.3 /opt/app/solr
准备 SolrCloud 配置
启动 SolrCloud 前,需将 Solr 核心配置上传到 Zookeeper 或启动时让 Solr 自动上传。本项目包含名为"product"的 collection,配置文件(schema.xml、solrconfig.xml、data-config.xml 等)已由开发准备。采取以下步骤:
- 配置 ZK_HOST:编辑 solr-7.7.3/bin/solr.in.sh,找到 ZK_HOST 设置,指向 Zookeeper 集群地址和 Solr 配置路径:
ZK_HOST="10.201.X.X:2181,10.201.X.X:2181,10.201.X.X:2181/solr_demo"
/solr_demo 表示 SolrCloud 将在 ZK 该路径下存储配置。确保该路径存在(之前用 zkCli 创建过)。
集群配置:可选在 solr.in.sh 中调整 JVM 堆大小(SOLR_HEAP)和垃圾回收策略。Solr 7 提供合理默认值;按需调整,如 SOLR_HEAP="1g"。
DataImportHandler 配置:若用 DIH 从 Cassandra 导入数据,需在 Solr 中加入 DIH 插件。Solr 7.7.3 提供 jar 包但默认未包在 solr-webapp 中。将以下 jar 复制到 solr-webapp/webapp/WEB-INF/lib/:
- dist/solr-dataimporthandler-7.7.3.jar
- dist/solr-dataimporthandler-extras-7.7.3.jar
- (如项目提供自定义 DIH 处理器如 solr-support-7.7.0-1.0.jar,也复制)
在一台机器上完成 Solr 准备,然后打包整个 solr 目录为 solr.zip 分发到各节点,确保配置一致。
上传 Solr 配置并启动 SolrCloud
假设开发提供的 Solr 核心配置(schema.xml、solrconfig.xml、data-config.xml 等)保存在 /home/netty/solr_core_config/product。将此上传到 Zookeeper:
$ cd /opt/app/solr
$ ./bin/solr zk upconfig -z 10.201.X.X:2181,10.201.X.X:2181,10.201.X.X:2181/solr_demo \
-n product -d /home/netty/solr_core_config/product
此命令将本地 product 配置目录上传到 ZK,注册为"product"配置集。(若需查看已存在的配置,可用 downconfig 下载,如 ./bin/solr zk downconfig -n product -d ./downloaded_conf -z ...)。
启动 Solr 节点
在每台 Solr 服务器以 SolrCloud 模式启动 Solr,加入集群:
$ ./bin/solr start -c -m 1g -z 10.201.X.X:2181,10.201.X.X:2181,10.201.X.X:2181/solr_demo -p 8983 -d /opt/app/solr/solr_data
参数说明:-c 表示 Cloud 模式;-m 1g 设最大堆 1GB;-z 指定 ZK 地址列表和路径;-p 8983 设端口(默认);-d /opt/app/solr/solr_data 指定 Solr 实例数据目录(将 example/server 内容复制于此以独立配置)。
启动所有 Solr 节点后,它们自动在 Zookeeper 注册并形成集群。接下来通过 Solr API 创建 collection。使用脚本为例:
$ ./bin/solr create -c product -n product -shards 2 -replicationFactor 2 -p 8983
创建名为"product"的 collection,使用已上传的"product"配置集,2 个分片、每片 2 个副本。脚本会在集群上创建 collection 并在各节点间分配 cores。成功后,SolrCloud 对外提供搜索服务。在浏览器访问任意 Solr 管理界面(如 http://10.201.X.X:8983/solr)查看集群状态和 collection 列表。
注意: SolrCloud 用 Zookeeper 存储配置和集群状态,务必保证 ZK 集群稳定。配置和启动 Solr 时,若 ZK 信息有误(地址或路径错误),Solr 将无法启动。
Python 微服务开发与部署
有了数据存储和索引组件,我们用 Python 实现微服务连接它们。原项目中 Java/Spring Boot 模块承担数据同步和搜索功能;我们用 Python 重写这些模块并提供示例。
微服务责任划分
数据同步服务(Kafka 消费者):持续消费 Kafka product_update 主题。商品数据更新消息到达时,服务用 DataStax cassandra-driver 写入 Cassandra,然后触发 Solr 索引更新。索引更新可增量(通过 Solr HTTP API)或批量(通过 DIH),取决于策略。Python 可用 kafka-python 或 confluent-kafka 消费消息,用 cassandra-driver 写 Cassandra,用 requests 或 pysolr 调 Solr API。
搜索 API 服务:暴露搜索 REST 接口。客户端请求由此服务接收并转发给 SolrCloud,格式化结果后返回。基于 Flask 或 FastAPI 实现。内部用 pysolr 或直接 HTTP 请求查询 Solr,将结果转为 JSON。此服务对应原 Java 搜索查询模块,也可从 Cassandra 补充信息——如从搜索索引得到产品 ID 列表,再从 Cassandra 查最新库存(如果这些信息未索引)。
辅助服务:如批量导入(product-upload)或定时任务触发(job-trigger),也可用 Python 实现。若需与既有调度框架(如 xxl-job)集成,可调其 HTTP 接口或按协议实现执行器。
微服务代码和启动脚本示例
用 FastAPI 构建搜索 API 服务并示范启动脚本。假设搜索服务代码为 app.py(简化示例):
# app.py (FastAPI 简易示例)
from fastapi import FastAPI
import pysolr
solr = pysolr.Solr('http://10.201.X.X:8983/solr/product', timeout=10) # Solr 地址
app = FastAPI()
@app.get("/search")
def search(q: str):
# 在 Solr 中查询
results = solr.search(q)
# 提取需要的字段返回
docs = [doc for doc in results]
return {"query": q, "results": docs}
生产环境可能需更复杂的查询构造和结果处理。接下来编写管理脚本 service_control.py 来启动或停止此 FastAPI 服务。用 subprocess 调 Uvicorn 运行 FastAPI 应用,模拟传统 Shell 中 java -jar ... & 的做法:
#!/usr/bin/env python3
# service_control.py
import subprocess, sys, os, signal
APP_COMMAND = ["uvicorn", "app:app", "--host", "0.0.0.0", "--port", "7220"]
def start():
"""启动服务"""
# 将输出重定向到日志文件
logfile = open("service.log", "a")
# 使用 nohup & 类似效果启动子进程
process = subprocess.Popen(APP_COMMAND, stdout=logfile, stderr=logfile, preexec_fn=os.setpgrp)
print(f"Service started with PID {process.pid}")
def stop():
"""停止服务"""
# 查找运行中的进程(通过端口或命令名)
try:
# 利用 pgrep 查找 uvicorn 进程
result = subprocess.run(["pgrep", "-f", "uvicorn.*7220"], capture_output=True, text=True)
pids = result.stdout.strip().split()
if not pids:
print("Service is not running.")
return
for pid in pids:
os.kill(int(pid), signal.SIGTERM)
print("Service stopped.")
except Exception as e:
print(f"Error stopping service: {e}")
def restart():
"""重启服务"""
stop()
start()
if __name__ == "__main__":
if len(sys.argv) < 2:
print("Usage: service_control.py [start|stop|restart]")
sys.exit(1)
cmd = sys.argv[1]
if cmd == "start":
start()
elif cmd == "stop":
stop()
elif cmd == "restart":
restart()
else:
print("Unknown command:", cmd)
此脚本模拟 Java 服务启动脚本逻辑:
- 用 uvicorn 运行 FastAPI 应用(监听 7220 端口,模拟原 Java 服务端口)。
- start() 用 subprocess.Popen 启动子进程并置于独立进程组(相当于 nohup)。日志输出到 service.log。
- stop() 用系统命令 pgrep 查找匹配 uvicorn 和端口的进程,发送 SIGTERM 关闭(或记录 process.pid 或用 psutil 库)。
- restart() 顺序调用 stop 再 start。
将 service_control.py 放服务器并赋予可执行权限,可用 ./service_control.py start 启动 Python 微服务,./service_control.py stop 停止服务。数据同步服务(Kafka 消费者)可采用类似方式启动,如编写 consumer.py 脚本启动 Kafka 消费循环,用控制脚本管理。注意: 生产环境建议用更健壮的方式托管微服务(如 Supervisor、systemd 或容器编排等),此处脚本仅作示范。
集群部署注意事项
所有组件部署完成后,还需关注以下集群搭建细节和最佳实践:
- 配置管理与一致性:集群配置繁多易错。使用配置管理工具(Ansible、Chef)或脚本确保机器配置一致。solr.in.sh、cassandra.yaml 等在不同节点需保持一致的参数应统一修改。自定义配置文件(如 Solr schema)确保版本正确、上传到 Zookeeper、在所有 Solr 节点生效。
- 资源分配:Cassandra、Solr、Kafka 对内存和 IO 敏感。根据机器规格调整 JVM 堆参数:
- Cassandra 默认使用一半物理内存作为堆;生产按数据规模调整但避免过大(长 GC 停顿)。
- Solr 堆依索引数据大小;保证常用查询工作集能容纳。若用排序/聚合,也要增加堆。索引放本地磁盘时,确保磁盘空间充足且 IO 性能好。
- Kafka 对文件系统顺序写友好,但日志目录须有足够空间,同时为 JVM 堆和操作系统页面缓存预留内存。Kafka 本身堆通常无需过大(几 GB 即可),更多依赖操作系统缓存加速磁盘 IO。
- 网络与端口:多节点集群部署时,防火墙需放行必要端口:
- Cassandra:9042(CQL)、7000/7001(集群通信)、7199(JMX)。
- Zookeeper:2181(客户端)、2888/3888(集群内部)需内网互通。
- Kafka:9092(或自定义 listeners 端口)需消费者/生产者访问;跨机房或容器环境需设 advertised.listeners 为可达地址。
- Solr:8983(默认)及节点间通信端口。防火墙隔离环境需开放 Solr 查询端口和 Zookeeper 端口。
- 数据导入与一致性:首次部署需将已有产品数据导入 Cassandra 和 Solr。可编写批处理脚本读源(CSV、数据库)写 Cassandra,再用 Solr DIH 全量导入或 Solr API 批量索引。之后增量数据走 Kafka -> 消费者服务自动更新。确保 Kafka 消费、Cassandra 写入、Solr 更新三者事务一致性:可先写 Cassandra、再写 Solr,Solr 更新失败可有补偿机制(如定期全量同步修复)。Cassandra 本身是最终一致性模型,写入成功即完成,无事务回滚,应用层要做好失败重试。
- 监控和日志:部署完成后,为各组件建立监控:
- Cassandra:监控节点延迟、读写吞吐、Compaction 状态等,及时扩容或调参。
- Solr:监控查询 QPS、索引大小、缓存命中率等,通过接口或 JMX 获取指标。
- Kafka:监控消息堆积(Lag)、消费者组状态、磁盘使用等。
- Python 微服务:记录关键操作日志(消费消息数、查询耗时等),用守护进程确保异常退出时自动重启。注册到 Supervisor 或 systemd 增强稳定性。
- 故障演练:生产前应测试各部分容错性。如停止一个 Cassandra 节点看查询是否受影响(需应用层考虑一致性级别),停止一个 Solr 节点检查查询是否自动切换副本,Kafka 任一 broker 挂掉后消息是否仍正常发送消费。验证通过后再正式上线。
总结与经验教训
通过上述实践,我们成功搭建了基于 Cassandra、SolrCloud、Zookeeper、Kafka 和 Python 微服务的分布式搜索系统。相比传统单体搜索应用,这种架构具有良好的扩展性和解耦性:Cassandra 提供高写入性能和水平扩展能力,SolrCloud 提供强大的搜索索引功能,Kafka 作为异步解耦的管道,Python 微服务让业务逻辑实现灵活高效。
从这次实践得到以下经验教训:
- 合理的架构设计:引入多种组件时,要明确各自职责,设计好数据流转路径。本案例中采用 Kafka 进行数据同步解耦,使系统更松耦合、可伸缩。若数据变化频率不高,也可直接由应用触发 Solr 更新,但引入消息队列提升长期灵活性。
- 充分的测试验证:集群配置繁杂,部署完成后需反复测试。曾遇到 Solr 节点无法加入集群,后发现是 ZK 路径配置错误;Kafka 消费延迟过高,追查发现消费者处理过慢。通过逐一定位问题并调整(修改配置或优化逻辑),最终使各部分协调工作。
- Python 替代 Java 的思考:用 Python 重构微服务模块,大幅减少了样板代码和编译部署时间,开发效率提升。但要注意 Python 在多线程、多进程方面的差异,充分利用异步 IO 或多进程。对性能要求极高的场景,需慎重评估 Python 成本(可考虑关键部分用 Cython 或 JNI 优化)。在我们实践中,Python 完全胜任数据消费和查询 API 工作,且易于维护。
- 配置与维护:集中管理配置是运维成功的关键。配置文件纳入版本控制,使用配置管理工具批量部署。敏感信息(密码、密钥)使用加密或配置中心管理。系统上线后要定期维护,包括 Cassandra 列压缩、Solr 索引优化、Kafka 日志清理等,保证长期稳定运行。
这次搭建与部署实践展示了较完整的搜索解决方案。从零开始组建这样一个系统需跨越多个技术领域知识,通过一步步安装配置和调优,我们掌握了各组件协同工作的要点。未来扩展中,可考虑将此架构容器化部署于 Kubernetes,实现更自动化的扩容和运维。希望这些经验对正在构建类似搜索平台的工程师有所帮助。