pathwaycom/pathway
Python ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.
Pathway – 实时数据框架
是什么 – Pathway 是一个开源 Python 库,用于构建可在批处理和流式数据上运行的数据管道。其底层基于 Differential Dataflow 的高性能 Rust 引擎执行 Python 定义的逻辑,让您在使用纯 Python 的同时,享受多线程、多进程和分布式执行的性能优势。
对 AI 为何重要 – 该框架自带 LLM x‑pack,包含包装器、解析器、嵌入器、分词器和内存中的向量索引,以及与 LangChain 和 LlamaIndex 的集成。这使得您可以轻松创建实时检索增强生成(RAG)或其他 LLM 驱动的工作流,对新数据即时响应。
核心功能
- 统一的批处理与流处理 – 一次编写管道,即可在静态文件或实时流上运行,无需修改。
- 丰富的连接器生态系统 – 内置支持 Kafka、PostgreSQL、GDrive、SharePoint、Airbyte(300+ 数据源),并支持自定义 Python 连接器。
- 状态化操作符 – Join、窗口化、排序和任意 Python UDF 在 Rust 引擎上运行,速度极快。
- 持久化与一致性 – 自动状态检查点;免费版提供 至少一次 语义,企业版支持 精确一次。
- 可扩展执行 – 本地多线程运行,或在 Docker/Kubernetes 上运行容器;企业版支持分布式云部署。
- LLM 辅助工具 – 提供主流 LLM 服务的即用型包装器、实时向量索引,兼容 LangChain/LlamaIndex。
典型应用场景
- 实时 ETL(例如:摄入 Kafka 主题,转换后加载到数据库)。
- 事件驱动的告警管道。
- 持续分析,如实时回归或仪表盘。
- 实时 LLM/RAG 应用,可即时摄入新文档、嵌入并回答查询。
快速上手
pip install -U pathway # macOS 或 Linux 上的 Python 3.10+ 支持
一个最小的“正数之和”管道示例:
import pathway as pw
class Input(pw.Schema):
value: int
src = pw.io.csv.read("./input/", schema=Input)
pos = src.filter(src.value >= 0)
out = pos.reduce(sum_val=pw.reducers.sum(pos.value))
pw.io.jsonlines.write(out, "output.jsonl")
pw.run()
通过 python my_pipeline.py 或 pathway spawn python my_pipeline.py 运行脚本,让引擎自动管理线程。
部署选项
- 本地 – 仅需导入
pathway并调用pw.run()。 - Docker – 使用官方
pathwaycom/pathway镜像,或在任意 Python 基础镜像中通过pip安装。 - Kubernetes / 云 – 企业版支持容器扩展;Render 上的快速启动文档已提供。
- 监控 – 内置仪表板可查看连接器吞吐量和延迟。
性能 – 基准测试表明,在需要时间连接或迭代算法的工作负载中,其吞吐量高于 Flink、Spark 和 Kafka Streams。
文档与支持 – 完整文档见 https://pathway.com/developers/,包含 API 参考、示例模板、Discord 社区和问题追踪器。
许可证 – Business Source License 1.1(非商业用途及大多数商业用途免费)。四年之后核心代码将重新授权为 Apache 2.0。支持仓库采用 MIT 许可证。
总结 – Pathway 让数据工程师和 AI 开发者能够用 Python 编写以“Rust 速度”运行的管道,同时处理历史数据和实时数据,并为以 LLM 为中心的 RAG 工作流提供一流支持。
相关
- 项目
- 项目
- 项目
- 项目
- 项目