用 Pathway 实时计算 K 部最佳电影:Kafka 与 Redpanda 双版本端到端实战
本指南基于仓库中的 best-movies-example 示例项目展开:它用 Pathway 实时数据框架实现了一个完整的端到端应用,在一个电影评分数据集上持续计算评分最高的 K 部电影(K-best-rated items),并演示了底层消息平台在 Kafka 与 Redpanda 之间切换时,应用代码几乎无需改动。读完本文,你将掌握示例的三层容器架构、从 CSV 静态文件构造实时流的做法、Pathway 中聚合与 Top-K 的表格变换写法,以及用 docker-compose 一键启动与结果查验的完整流程。
示例的目标与整体设计
顶层说明文档开宗明义地给出这个项目想要证明的两件事:
- 使用 Pathway 实时数据框架,对一个电影数据集完成 K 部最佳评分电影 的端到端实时计算;
- 演示 从 Kafka 平滑切换到 Redpanda 的过程。
其核心意图是证明:无论底层选择 Kafka 还是 Redpanda,Pathway 实时数据框架都能以完全相同的方式被使用。因此仓库将该示例平行地复制为两个自包含的子项目:
- examples/projects/best-movies-example/kafka-version:基于 Kafka/ZooKeeper 的消息平台版本;
- examples/projects/best-movies-example/redpanda-version:基于 Redpanda 的版本。
两个子项目内部布局完全镜像,都由三类组件构成:
| 组件 | 职责 |
|---|---|
| 流式消息平台 | Kafka/ZooKeeper(Kafka 版)或 Redpanda(Redpanda 版),承担消息中转 |
| Producer 容器 | 读取静态 CSV 文件,把每行评分记录序列化为 JSON 消息,逐条写入 Kafka/Redpanda 的 ratings topic,从而构造出一条“流” |
| Pathway 容器 | 订阅 ratings topic,用 Pathway 实时数据框架持续计算评分最高的 K 部电影,并把结果写入 CSV 文件 |
所有容器通过 docker-compose 统一编排管理,每个子项目目录下都附带一个 Makefile 封装常用操作命令。
目录与文件布局
整个示例在仓库中的物理结构如下(以 Kafka 版为例,Redpanda 版完全对称):
examples/projects/best-movies-example/
├── README.md # 顶层项目说明
├── kafka-version/
│ ├── README.md # Kafka 版使用说明
│ ├── docker-compose.yml # 三容器编排定义
│ ├── Makefile # 常用命令封装
│ ├── pathway-src/
│ │ ├── Dockerfile # Pathway 计算容器镜像定义
│ │ └── process-stream.py # Pathway 消费 + Top-K 计算逻辑
│ └── producer-src/
│ ├── Dockerfile # Producer 容器镜像定义
│ ├── create-stream.py # CSV → Kafka 的流生成脚本
│ └── dataset.csv # 电影评分玩具数据集
└── redpanda-version/
└── ... # 与 kafka-version 结构一致
值得注意的是 kafka-version 与 redpanda-version 两套源码的唯一实质差异集中在 消息平台地址 与 Pathway 读连接器的选择 上,其余(数据处理逻辑、容器编排骨架、启动流程)几乎逐字相同——这正是该示例想要强调的“与底层消息系统解耦”。
输入数据:按 MovieLens 格式构造的评分流
Producer 容器内提供了一个玩具数据集 dataset.csv,其格式与公开的 MovieLens25M 数据集保持一致,采用四列逗号分隔,首行为表头:
userId,movieId,rating,timestamp
1,296,5.0,1147880044
1,306,3.5,1147868817
1,307,5.0,1147868828
1,665,5.0,1147878820
四列的语义分别为:用户 ID、电影 ID、评分值(0~5 浮点)、评分时间戳(Unix 秒)。项目说明特别提到这是一个“玩具”规模的样本,仅用于演示与验证——如果希望测试完整规模,作者鼓励使用整个 MovieLens25M 数据集自行替换。
流平台的两种编排方式对比
Kafka 版:Kafka + ZooKeeper
docker-compose.yml 定义了四个服务,其中消息平台由 zookeeper 与 kafka 两个容器组成:
zookeeper:使用confluentinc/cp-zookeeper:5.5.3,监听 2181 端口;kafka:使用confluentinc/cp-enterprise-kafka:5.5.3,depends_onzookeeper。环境变量中关键配置包括:KAFKA_AUTO_CREATE_TOPICS=true(允许 topic 自动创建)、KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092(容器内互访使用主机名kafka)、KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1(单节点部署)等;- 容器启动命令在等待 15 秒后,后台执行
kafka-topics --create --zookeeper zookeeper:2181 --replication-factor 1 --partitions 1 --topic ratings,预创建单分区 topicratings,随后才拉起 broker 主进程; pathway与stream-producer两个业务容器均depends_onkafka,并接入同一个 bridge 网络best-movies-kafka-network,从而能够通过主机名kafka直接访问 broker。
Redpanda 版:一个进程替代整套 Kafka 栈
Redpanda 版 docker-compose.yml 用单个 redpanda 服务(镜像 docker.redpanda.com/redpandadata/redpanda:v23.1.2,--mode dev-container 开发模式)同时替代了 ZooKeeper + Kafka 两个组件。其启动参数值得关注:
- 通过
--kafka-addr/--advertise-kafka-addr声明internal://redpanda:9092供容器网络内部使用,external://localhost:19092供宿主机映射访问,从而保持 Kafka 协议端口 9092 这一对外语义不变; - 开启
redpanda.enable_transactions=true、redpanda.enable_idempotence=true(幂等与事务支持)以及redpanda.auto_create_topics_enabled=true(自动建 topic),因此无需像 Kafka 版那样手工预创建ratingstopic; - 数据卷
redpanda:/var/lib/redpanda/data用于持久化 broker 数据。
从编排层就能看出切换成本之低:Producer 与 Pathway 容器的镜像构建、启动方式完全相同,唯一变化是消息平台服务名从 kafka 变成了 redpanda(两者在容器网络内都以 :9092 暴露 Kafka 协议端点)。
容器镜像与运行环境
两个业务容器的镜像定义非常精简:
- pathway-src/Dockerfile 基于
python:3.10,执行pip install -U pathway安装最新版 Pathway SDK,拷入process-stream.py后以python -u process-stream.py启动(-u保证日志实时输出); - producer-src/Dockerfile 同样基于
python:3.10,安装kafka-python库,拷入脚本与数据集后启动create-stream.py。
需要留意的是,两份 process-stream.py 顶部都包含一行 License 设置(详见下文“License 与版本说明”一节)。
数据生产端:把静态 CSV 变成实时消息流
create-stream.py 扮演“模拟实时数据源”的角色,其流程是:
time.sleep(30)先等待消息平台完全就绪;- 用
kafka-python的KafkaProducer连接kafka:9092(协议PLAINTEXT,兼容 API 版本 0.10.2); - 逐行读取
dataset.csv,跳过表头后,把userId / movieId / rating / timestamp组装成 JSON 字典发送到ratingstopic; - 每发送一条记录后
time.sleep(0.1)——人为插入 100ms 间隔,以模拟低吞吐的持续到达流,而不是一次性灌入; - 数据发送完毕后,向 topic 写入一条文本消息
*COMMIT*作为批次的收尾同步标记。
Redpanda 版脚本 redpanda-version/producer-src/create-stream.py 逻辑一致,仅把 bootstrap_servers 改为 ["redpanda:9092"],并把 *COMMIT* 标记调整到数据发送之前发出。注意:该标记是一条字符串消息,并非评分数据记录,示例通过它给出流的阶段性边界信号。
Pathway 消费端:从读入到 Top-K 输出的全链路
Pathway 计算容器的主程序 process-stream.py 完整地呈现了“接入 → 增量聚合 → Top-K → 落盘”的实时计算管线,全文核心如下:
import time
import pathway as pw
pw.set_license_key("demo-license-key-with-telemetry")
K = 3
rdkafka_settings = {
"bootstrap.servers": "kafka:9092",
"security.protocol": "plaintext",
"group.id": "0",
"session.timeout.ms": "6000",
}
def compute_best(t_ratings, K):
# ① 按电影聚合:总评分 + 评分数
t_best_ratings = t_ratings.groupby(pw.this.movieId).reduce(
pw.this.movieId,
sum_ratings=pw.reducers.sum(pw.this.rating),
number_ratings=pw.reducers.count(),
)
# ② 计算平均分(防除零)
t_best_ratings = t_best_ratings.select(
pw.this.movieId,
pw.this.number_ratings,
average_rating=pw.apply(
lambda x, y: (x / y) if y != 0 else 0,
pw.this.sum_ratings,
pw.this.number_ratings,
),
)
# ③ 组装用于排序的比较元组 (平均分, 评分数, movieId)
t_best_ratings = t_best_ratings.select(
movie_tuple=pw.apply(
lambda x, y, z: (x, y, z),
pw.this.average_rating,
pw.this.number_ratings,
pw.this.movieId,
)
)
# ④ 全局归约为一个升序有序元组
t_best_ratings = t_best_ratings.reduce(
total_tuple=pw.reducers.sorted_tuple(pw.this.movie_tuple)
)
# ⑤ 取排序后最末 K 个 → Top-K
t_best_ratings = t_best_ratings.select(
K_best=pw.apply(lambda my_tuple: (list(my_tuple))[-K:], pw.this.total_tuple)
)
# ⑥ 展开为逐行输出
t_best_ratings = t_best_ratings.flatten(t_best_ratings.K_best).select(
pw.this.K_best
)
# ⑦ 还原为命名列
t_best_ratings = t_best_ratings.select(
movieId=pw.apply(lambda rating_tuple: rating_tuple[2], pw.this.K_best),
average_rating=pw.apply(lambda rating_tuple: rating_tuple[0], pw.this.K_best),
views=pw.apply(lambda rating_tuple: rating_tuple[1], pw.this.K_best),
)
return t_best_ratings
class inputStreamSchema(pw.Schema):
movieId: int
rating: float
t_ratings = pw.io.kafka.read(
rdkafka_settings,
topic="ratings",
format="json",
schema=inputStreamSchema,
autocommit_duration_ms=100,
)
t_best_ratings = compute_best(t_ratings, 3)
pw.io.csv.write(t_best_ratings, "./best_ratings.csv")
time.sleep(20) # 等待 Kafka 就绪
pw.run()
逐步拆解 compute_best 的七段变换
上述函数是一段颇具代表性的 Pathway 声明式数据流代码,逐段理解即可掌握“流式 Top-K”的通用范式:
- 分组聚合:
groupby(pw.this.movieId).reduce(...)按电影 ID 分组,用pw.reducers.sum与pw.reducers.count维护每个电影累计的评分总和与评分条数——这是后续计算平均分的基础(pw.reducers是 Pathway 内置的增量聚合算子集合); - 计算平均分:
pw.apply(lambda x, y: ...)把任意 Python 函数逐行应用到列上,此处用sum/count求均值并防御y == 0的除零情况; - 构造比较元组:把
(average_rating, number_ratings, movieId)打包成一个 Python 元组。由于元组天然按元素顺序比较,后续排序会先按平均分、再按评分数、最后按 movieId 决定次序; - 全局归约排序:
reduce(pw.reducers.sorted_tuple(...))把所有电影的元组收敛到一张单行表,列中保存的是按上述规则升序排列的完整元组序列; - 截取 Top-K:对升序序列执行
list(my_tuple)[-K:]取最末 K 个元素,即平均分最高的 K 部电影(平均分相同时,评分数更多、movieId 更大的排在更后、更易入选); - 展平:
flatten把单行中的 K 元组展开为 K 行独立的输出; - 列还原:把三元组按下标
[0]/[1]/[2]分别还原为average_rating、views(即评分数)、movieId三个具名列。
输入 schema 与读连接器配置
输入结构用 pw.Schema 子类声明:
class inputStreamSchema(pw.Schema):
movieId: int
rating: float
尽管 Producer 发送的 JSON 携带四个字段,这里只声明了计算真正需要的 movieId 与 rating,其余字段会被 JSON 解析层忽略。
读取端参数中两个细节值得展开:
rdkafka_settings是一份标准的 librdkafka 风格连接配置(bootstrap.servers、security.protocol、group.id、session.timeout.ms),Pathway 的 Kafka/Redpanda 连接器都直接消费这份字典;autocommit_duration_ms=100控制消费偏移量(offset)的自动提交周期为 100ms,保证计算进度持续持久化、可恢复。
Kafka 版与 Redpanda 版的“唯一差异”
对比两份 process-stream.py,Redpanda 版 redpanda-version/pathway-src/process-stream.py 只做了两处改动:
rdkafka_settings["bootstrap.servers"]由"kafka:9092"改为"redpanda:9092";- 数据接入调用由
pw.io.kafka.read(...)改为pw.io.redpanda.read(...)。
其余——包括 compute_best 函数、schema、输出写 CSV、pw.run() 启动方式——完全一致。这份“同构性”并非偶然:查看 python/pathway/io/redpanda/init.py 的源码可以看到,pw.io.redpanda.read 在实现上正是转调了 kafka.read(内部走同一条 librdkafka 消费链路),仅以 Redpanda 的兼容协议为前置假设。这从源码层面印证了 README 的核心主张:无论 Kafka 还是 Redpanda,只要对外暴露 Kafka 0.10+ 兼容协议,Pathway 应用代码就无需为平台差异而分叉。
输出与运行
计算得到的最终表通过 pw.io.csv.write(t_best_ratings, "./best_ratings.csv") 持续写入容器内的 best_ratings.csv,每行对应一部入选电影的 movieId、average_rating 与 views。随后脚本 time.sleep(20) 给消息平台留出启动窗口,最后调用 pw.run() 阻塞式启动实时计算引擎,此后所有新增评分都会自动触发结果的增量更新。
一键启动与结果查验
两个版本的启动方式完全相同,README 给出的标准操作如下。
第一步:拉起整套应用
docker compose up -d
或者直接使用封装好的 Makefile 目标 make build(其内部执行的正是 docker compose up -d)。-d 让三个/四个容器在后台运行,无需占用前台终端。
第二步:进入 Pathway 容器查看结果
make connect
# 等价于:docker compose exec -it pathway bash
进入容器后打印结果文件:
cat best_ratings.csv
即可看到当前评分最高的 K(此处 K=3)部电影及其平均分、评分数。由于应用是实时增量计算,随着 Producer 以 100ms 间隔持续灌入新评分,该文件内容会不断被 Pathway 更新。
第三步(可选):查看各容器日志
各版本 Makefile 中都提供了按组件细分的日志查看目标。Kafka 版 Makefile 提供的目标整理如下:
| Make 目标 | 作用 | 底层命令 |
|---|---|---|
build |
构建并后台启动全部容器 | docker compose up -d |
stop |
停止并清理(含删除数据卷与本地镜像) | docker compose down -v + docker rmi ... |
connect |
进入 Pathway 计算容器 | docker compose exec -it pathway bash |
connect-prod |
进入 Producer 容器 | docker compose exec -it stream-producer bash |
connect-kafka |
进入 Kafka 容器 | docker compose exec kafka bash |
logs |
查看 Pathway 容器日志 | docker compose logs pathway |
logs-prod |
查看 Producer 日志 | docker compose logs stream-producer |
logs-kafka / logs-zookeeper |
查看消息平台日志 | docker compose logs kafka/zookeeper |
Redpanda 版 Makefile 除上述对应目标外,还额外提供了基于 Redpanda 自带运维工具 rpk 的几个目标:info(rpk cluster info 查看集群信息)、create-topic、create-message。需要提醒的是:该 Makefile 中 connect-redpanda、info 等目标引用的容器名是 redpanda-tuto,而当前 docker-compose 中容器实际命名为 redpanda(见 redpanda-version/docker-compose.yml),两者存在命名不一致;如需使用这些目标,可将其中的容器名替换为 redpanda,或改用 docker compose exec redpanda <cmd> 的形式。
运行节奏的协调机制
由于容器间存在启动依赖,两版代码都通过显式 time.sleep 做了时序兜底,理解这些延时有助于排查“启动后迟迟无结果”的问题:
- Kafka 版中,broker 容器启动后先睡 15 秒再建 topic;Producer 启动后睡 30 秒才开始发送;Pathway 脚本在
pw.run()前也睡 20 秒等待消息平台就绪; - Redpanda 版中,Redpanda 单容器启动较快且自动建 topic,Producer 与 Pathway 仍保留了类似的等待窗口以保证
redpanda:9092已可连接。
这套“sleep 兜底 + depends_on”的组合在真实生产环境通常会被健康检查(healthcheck)取代,但对一个教学型示例而言,它保证了在任意机器上 docker compose up -d 后各组件能以正确的相对次序完成初始化。
License 与版本说明
两份 process-stream.py 顶部都有如下注释与调用:
# To use advanced features with Pathway Scale, get your free license key ...
# To use Pathway Community, comment out the line below.
pw.set_license_key("demo-license-key-with-telemetry")
按其说明:Pathway 提供 Community 与面向高级特性的商业版本两种形态。示例默认以带遥测的演示 License Key 运行;如果你希望以 Community 形态运行,直接注释掉 pw.set_license_key(...) 这一行即可。此外镜像采用 pip install -U pathway 安装最新发布版本,因此运行结果会跟随 Pathway 的最新发行版行为。
延伸:如何迁移到更大的数据集
示例刻意把输入数据与业务逻辑解耦,这让替换数据集非常直接:只要目标 CSV 保持 userId,movieId,rating,timestamp 四列(无表头行需删除或调整 Producer 中跳过首行的逻辑),把文件替换为完整版 MovieLens25M 数据后重新 docker compose build 即可。由于 Pathway 是增量流式计算,无论数据规模多大,其内部维护的 sum/count 聚合状态与 Top-K 排序都会随每条新消息即时更新,这正是本示例希望传达的“实时数据处理”核心体验。
小结
best-movies-example 是一个精心设计的对照实验:同一份数据处理需求,在 Kafka 与 Redpanda 两套消息平台下各跑一遍,借此证明 Pathway 对 Kafka 协议生态的兼容设计——从 python/pathway/io/redpanda/init.py 中 redpanda.read 直接转调 kafka.read 的实现可知,这种“换平台不改业务代码”的能力是框架层面的架构决策,而非示例特例。对初学者而言,compute_best 中“groupby/reduce 聚合 → 构造元组 → sorted_tuple → 截尾取 Top-K → flatten”的组合,也是一份可直接复用到排行榜、热榜、实时榜单类场景的流式 Top-K 参考实现。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0627
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00