文章 · 2022-02-14

基于 Cassandra、SolrCloud、Zookeeper、Kafka 和 Python 微服务的搜索系统搭建实践

本方案假设 Linux 服务器环境(CentOS 或 Ubuntu)已安装适当版本的 Java JDK(Cassandra、Solr、Kafka 都依赖 Java)。多台服务器部署方案如下:

各组件通过内部网络通信,防火墙需放行必要端口(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,按网络环境调整:

seed_provider:
    - class_name: org.apache.cassandra.locator.SimpleSeedProvider
      parameters:
          - seeds: "10.201.X.X,10.201.X.X"

至少包含集群中一个节点 IP 作为种子。

启动 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 下创建多张表:

搜索主要由 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,关注以下配置:

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 为选举端口。

启动 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:

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 等路径。

根据需要调整内存等参数,但默认配置适合初始测试。

启动 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="10.201.X.X:2181,10.201.X.X:2181,10.201.X.X:2181/solr_demo"

/solr_demo 表示 SolrCloud 将在 ZK 该路径下存储配置。确保该路径存在(之前用 zkCli 创建过)。

在一台机器上完成 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 重写这些模块并提供示例。

微服务责任划分

  1. 数据同步服务(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。

  2. 搜索 API 服务:暴露搜索 REST 接口。客户端请求由此服务接收并转发给 SolrCloud,格式化结果后返回。基于 Flask 或 FastAPI 实现。内部用 pysolr 或直接 HTTP 请求查询 Solr,将结果转为 JSON。此服务对应原 Java 搜索查询模块,也可从 Cassandra 补充信息——如从搜索索引得到产品 ID 列表,再从 Cassandra 查最新库存(如果这些信息未索引)。

  3. 辅助服务:如批量导入(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 服务启动脚本逻辑:

将 service_control.py 放服务器并赋予可执行权限,可用 ./service_control.py start 启动 Python 微服务,./service_control.py stop 停止服务。数据同步服务(Kafka 消费者)可采用类似方式启动,如编写 consumer.py 脚本启动 Kafka 消费循环,用控制脚本管理。注意: 生产环境建议用更健壮的方式托管微服务(如 Supervisor、systemd 或容器编排等),此处脚本仅作示范。

集群部署注意事项

所有组件部署完成后,还需关注以下集群搭建细节和最佳实践:

总结与经验教训

通过上述实践,我们成功搭建了基于 Cassandra、SolrCloud、Zookeeper、Kafka 和 Python 微服务的分布式搜索系统。相比传统单体搜索应用,这种架构具有良好的扩展性和解耦性:Cassandra 提供高写入性能和水平扩展能力,SolrCloud 提供强大的搜索索引功能,Kafka 作为异步解耦的管道,Python 微服务让业务逻辑实现灵活高效。

从这次实践得到以下经验教训:

这次搭建与部署实践展示了较完整的搜索解决方案。从零开始组建这样一个系统需跨越多个技术领域知识,通过一步步安装配置和调优,我们掌握了各组件协同工作的要点。未来扩展中,可考虑将此架构容器化部署于 Kubernetes,实现更自动化的扩容和运维。希望这些经验对正在构建类似搜索平台的工程师有所帮助。

© 2026 Yuxu Ge ·