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.pypathway 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 工作流提供一流支持。

相关

  • 项目
  • 项目
  • 项目
  • 项目
  • 项目