尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Pathway 实战:用 Kafka/Redpanda 流式计算电影数据集 Top-K 评分(best-movies-example 模板全解)

Pathway 实战:用 Kafka/Redpanda 流式计算电影数据集 Top-K 评分(best-movies-example 模板全解) Pathway 实战用 Kafka/Redpanda 流式计算电影数据集 Top-K 评分best-movies-example 模板全解【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway本文围绕仓库中 docs/2.developers/7.templates/ETL/_readmes/best-movies-example.md 所对应的示例项目展开它是一套基于 Pathway 的端到端流处理模板用于在 MovieLens 格式的评分数据集上实时计算K 部评分最高的电影。读完本文你将掌握该模板的三层容器架构流平台 生产者 Pathway 计算、Kafka 与 Redpanda 两种部署版本的无缝切换方法以及如何用 docker-compose 一键启动并用 Makefile 完成查看结果、查看日志等日常操作。一、模板要解决的两个问题该模板项目同时见 examples/projects/best-movies-example/README.md的目标是双重的用 Pathway 构建一个端到端应用在电影数据集中计算 K-bestK 佳评分电影条目演示如何把底层流平台从 Kafka 切换到 Redpanda。其核心立意在于无论底层选择 Kafka 还是 RedpandaPathway 的应用代码可以以完全相同的方式运行。对于希望评估消息中间件可替换性的团队来说这是一个非常直观的 A/B 对照样本——你只需要替换 docker-compose 中定义的流平台服务Pathway 侧的数据接入、变换逻辑、结果输出均无需改动。二、项目结构与三种角色整个项目由三部分容器组成全部通过 docker-compose 统一编排角色职责说明流平台Kafka/ZooKeeper 或 Redpanda负责消息的接收与分发stream-producerPython 容器读取静态 CSV 文件逐条模拟产生流式消息发送到流平台pathwayPython 容器使用 Pathway 消费数据流实时计算 K-best 评分电影模板在仓库中同时提供了两套独立可运行的同构实现examples/projects/best-movies-example/kafka-version/基于 Kafka/ZooKeeper入口 README 见 kafka-version/README.mdexamples/projects/best-movies-example/redpanda-version/基于 Redpanda入口 README 见 redpanda-version/README.md。两个版本内部均采用相同的目录约定version/ # kafka-version 或 redpanda-version ├── docker-compose.yml # 服务编排 ├── Makefile # 常用操作命令封装 ├── pathway-src/ # Pathway 计算容器 │ ├── Dockerfile │ └── process-stream.py # 核心流处理逻辑计算 K-best └── producer-src/ # 数据生产者容器 ├── Dockerfile ├── create-stream.py # 把 CSV 变成消息流 └── dataset.csv # MovieLens 格式玩具数据集三、数据集与消息格式模板自带一个玩具数据集 producer-src/dataset.csv其字段格式与MovieLens 25M数据集保持一致userId, movieId, rating, timestamp四列。官方鼓励使用方在完整 MovieLens25M 数据上自行测试——只需把该 CSV 换成完整数据集即可验证真实规模下的处理能力。从生产者源码 create-stream.py 可以看到每行 CSV 会被转换为如下 JSON 消息写入 topic{userId: 1, movieId: 110, rating: 4.0, timestamp: 1425941529}而 Pathway 侧只消费其中的两个字段。在 process-stream.py 中通过 Schema 声明class inputStreamSchema(pw.Schema): movieId: int rating: float即整个计算仅依赖movieId与rating用户 ID 与时间戳会被自动忽略。四、运行方式docker compose 与 Makefile在两个版本各自的目录下启动服务即可命令完全一致# 一键构建并后台启动全部容器 docker compose up -d # 等价方式直接调用 Makefile make build等待容器启动后进入 Pathway 计算容器查看结果# 进入 pathway 容器两种写法等价 make connect # 或docker compose exec -it pathway bash # 打印计算结果文件 cat best_ratings.csv计算结果由 Pathway 持续写入./best_ratings.csv。由于容器内挂载的是容器本地文件系统make connect进入后cat best_ratings.csv即可看到实时更新的 Top-K 电影含movieId、average_rating、views三个字段。Makefile 常用目标一览两个版本的 Makefile 目标几乎对称这里以 kafka-version/Makefile 为例目标作用builddocker compose up -d一键后台启动stopdocker compose down -v并清理两个自建镜像pathway、stream-producerconnect进入 pathway 容器connect-prod进入 stream-producer生产者容器connect-kafka进入 kafka 容器Redpanda 版为connect-redpandalogs/logs-prod/logs-kafka/logs-zookeeper分别查看 pathway、生产者、Kafka、ZooKeeper 的日志Redpanda 版 Makefileredpanda-version/Makefile还额外提供了基于rpkRedpanda 官方 CLI的运维目标info集群信息、create-topic、create-message方便直接向 Redpanda 发送测试消息验证链路。五、底层原理剖析5.1 Pathway 端声明式拓扑 实时刷新Pathway 端逻辑集中在 kafka-version/pathway-src/process-stream.pyRedpanda 版本逻辑与其镜像对齐。它首先读取 Kafka topict_ratings pw.io.kafka.read( rdkafka_settings, topicratings, formatjson, schemainputStreamSchema, autocommit_duration_ms100, )其中rdkafka_settings为底层 librdkafka 风格的连接配置源码 L14-L19rdkafka_settings { bootstrap.servers: kafka:9092, # 容器网络内通过服务名访问 security.protocol: plaintext, group.id: 0, session.timeout.ms: 6000, }Kafka 版使用bootstrap.servers: kafka:9092broker 服务名对应 docker-compose 中 Kafka 的KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092Redpanda 版则指向redpanda:9092这正是两版在平台相关层唯一的本质差异点。K-best 的计算在compute_best函数中分五步完成源码 L22-L59按电影聚合groupby(pw.this.movieId).reduce(...)计算每部电影的评分总和sum_ratings与评分次数number_ratings计算平均分通过pw.apply应用lambda x, y: (x / y) if y ! 0 else 0得到average_rating对除零做了兜底构造复合元组把(average_rating, number_ratings, movieId)打包成一个元组便于整体排序全局有序聚合reduce(total_tuplepw.reducers.sorted_tuple(...))把所有电影按平均分排序聚合到单个元组中再用 Python 切片list(my_tuple)[-K:]取出最高的 K 个切片取自升序元组末尾故为 Top-K拆回列flatten(...).select(...)把每个结果元组还原为movieId、average_rating、views评分人数三列。最后结果输出到 CSV 并启动引擎pw.io.csv.write(t_best_ratings, ./best_ratings.csv) time.sleep(20) # 等待 Kafka/Redpanda 就绪 pw.run()需要强调的是这里的pw.run()一旦启动上述所有变换都是增量式实时维护的——每当新评分消息到达best_ratings.csv中的 Top-K 就会自动重算刷新这正是 Pathway 表引擎对流的连续查询语义在实时排行榜场景中的直接体现。5.2 生产者端把静态 CSV 模拟成实时流生产者 create-stream.py 使用kafka-python的KafkaProducer先time.sleep(30)等待 broker 就绪与服务依赖depends_on配合进一步规避启动竞态跳过 CSV 表头后逐行读取组装成 JSON 消息producer.send(topic, ...)发送每条消息之间time.sleep(0.1)人为限速形成节奏可控的流式回放效果发送完毕后写入一条特殊的*COMMIT*终止标记并关闭连接Pathway 端可据此感知数据回放完成。其运行环境由 producer-src/Dockerfile 定义基于python:3.10仅安装kafka-python随后拷贝脚本与数据集执行。5.3 平台差异只发生在 docker-compose 层对比两个 docker-compose 文件可以更清楚地看到切换平台的成本有多低Kafka 版docker-compose.yml定义zookeepercp-zookeeper:5.5.3客户端端口 2181与kafkacp-enterprise-kafka:5.5.3两个服务Kafka 启动命令还会在后台sleep 15后自动创建名为ratings的 topic1 分区、1 副本并开启KAFKA_AUTO_CREATE_TOPICS、设置KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092以打通容器网络内的消费。Redpanda 版docker-compose.yml用单个redpanda服务redpandadata/redpanda:v23.1.2同时承担 broker 与 schema registry 等职责无需 ZooKeeper启动参数中开启redpanda.auto_create_topics_enabledtrue从而同样支持 topic 自动创建其容器内 Kafka 兼容监听地址为redpanda:9092。而两个版本中的pathway与stream-producer服务定义几乎完全一致depends_on等待平台服务共享同一张 bridge 网络。也就是说切换本质上是换掉流平台那一组服务Pathway 与生产者逻辑保持不变。这与官方连接器相关文档 80.switching-to-redpanda.md 中Pathway 统一使用 rdkafka 参数访问 Kafka 协议兼容中间件的说明相互印证。六、从模板走向自己的实时 Top-K 应用要把这套模板改造成自己的场景只需沿三个触点调整换数据用你自己的movieId/rating或其他打分实体/分数CSV 替换 dataset.csv若字段含义变化则同步调整create-stream.py中 JSON 字段名与process-stream.py里的inputStreamSchema调 K 值修改 process-stream.py 顶部 的K 3也可改为由环境变量传入换输出把pw.io.csv.write换成任意 Pathway 输出连接器如数据库、其他消息系统等即可让 Top-K 结果实时流向你的下游。若希望同时对齐两个平台版本建议以 kafka-version 为基线实现并验证算法再复制一份目录、按 redpanda-version 的 docker-compose 替换流平台服务并将连接地址改为redpanda:9092即可在几分钟内完成另一套部署。七、小结best-movies-example 是一份小而完整的 Pathway 端到端模板它用最短的代码闭环展示了数据接入Kafka/Redpanda 读→ 声明式流变换增量聚合 Top-K→ 结果输出CSV 写的全过程并通过 Kafka/Redpanda 双版本证明了 Pathway 应用与底层消息平台之间的解耦。无论你是刚接触 Pathway 想跑通第一个实时任务还是在为生产环境评估 Kafka 与 Redpanda 的可替换性都可以直接以本模板作为起点相关模板综述另见 docs/2.developers/7.templates/ETL/_readmes/kafka-version.md 与 redpanda-version.md。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表