ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Pathway Live Data Framework 入门全景:从 Schema、连接器到 Rust 引擎的实时数据管道

Pathway Live Data Framework 入门全景:从 Schema、连接器到 Rust 引擎的实时数据管道 Pathway Live Data Framework 入门全景从 Schema、连接器到 Rust 引擎的实时数据管道【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway本文基于 Pathway 官方文档 Live Data Framework 概览 展开系统介绍如何用 Pathwaypip install pathway即可安装的 Python 流处理框架搭建一条完整的实时数据管道定义数据 Schema、用连接器接入 CSV/Kafka/SQLite 等数据源、执行过滤与聚合等转换、编写时间窗口与时态 Join 等时序算子最后配置输出连接器并通过pw.run()让 Rust 引擎持续运行。读完后你将掌握 Pathway 静态 Schema 动态表内容 增量计算 的核心编程模型并能写出可运行、可扩展的流式处理应用。1. 安装与导入一条命令接入框架Pathway 的编程层是纯 Python 的安装非常简单pip install pathway安装完成后像导入任意 Python 库一样导入它import pathway as pw从 python/pathway/init.py 的导出列表可以看到pw命名空间一次性暴露了框架的全部核心构件Schema、Table、LiveTable、apply、sql、groupby、join、run、this列引用语法、debug、io、udfs、demo等。也就是说一篇典型的 Pathway 应用代码只需要import pathway as pw这一行导入就能完成数据接入 → 转换 → 输出 → 运行的全部工作。框架的计算引擎则以 Rust 二进制形式随包分发见 python/pathway/_engine_finder.py 负责定位引擎这正是文档反复强调转换在底层用 Rust 编写因此非常高效的原因。2. 定义数据 Schema静态结构保证类型安全Pathway 中所有数据表都有 Schema用于声明列名与列类型保证数据组织有序、类型一致class InputSchema(pw.Schema): colA: int colB: float colC: str这个InputSchema声明了三列colA整数、colB浮点数、colC字符串。Schema 的作用有两层一是类型安全转换操作在运行前就能按声明类型校验二是运行时优化Rust 引擎可以基于已知类型做列式存储与向量化计算。支持的类型包括基础类型详见 数据类型文档bool、str、bytes、int、float更复杂的类型Optional可选类型、时间类型如datetime.datetime内部对应DateTimeNaive/DateTimeUtc/Duration等可从 python/pathway/init.py 的导出中确认。Schema 的底层实现位于 python/pathway/internals/schema.py。该文件末尾的Schema类通过元类SchemaMetaclasspython/pathway/internals/schema.py在子类被定义时收集列声明__init_subclass__还支持append_only、id_dtype等元参数来控制表的主键与追加语义——这些是理解后文表是动态内容 静态结构模型的关键主键id决定了一行数据在流中是新增/更新/删除中的哪一种。3. Tables内容动态、结构静态的数据容器Table 是 Pathway 中真正承载数据的对象它由若干列组成每列存放同类型数据组织方式与关系型数据库的表相似。与静态数据库表不同Pathway 的表是数据流的快照——Schema 固定不变内容则随新事件到达而实时更新。核心概念文档中对此有完整论述见 Core Concepts其中对新增、更新、删除三类事件如何反映为表行变化的解释值得细读。在源码层面表类型定义在 python/pathway/internals/table.py例如Table.filterpython/pathway/internals/table.py#L497与Table.groupbypython/pathway/internals/table.py#L1192都是该类的实例方法。文档概览示例中的pw.this.colA语法就是python/pathway/internals/column_namespace.py中列引用机制的体现让表达式写法接近 SQL 而不是裸函数调用。4. 输入连接器把数据源接成实时表要创建表需要连接器connector。连接器负责实时从外部数据源读取并摄取数据。以 CSV 为例input_table pw.io.csv.read(./data/, schemaInputSchema)文档概览中给出的常用输入连接器一览输入连接器示例CSVpw.io.csv.read(./data/, schemaInputSchema)Kafkapw.io.kafka.read(rdkafka_settings, topicexample, schemaInputSchema, formatcsv)SQLitepw.io.sqlite.read(./data_path/, table_name, schemaInputSchema)Google Drivepw.io.gdrive.read(object_id***, service_user_credentials_filecredentials.json)这些连接器在源码中都有对应的独立子包位于 python/pathway/io/ 目录csv、kafka、sqlite、gdrive、debezium、postgres、kinesis、mqtt、elasticsearch、iceberg、deltalake、s3、airbyte等数十个包规模远超概览表格所列。文档还特别提到 Airbyte 连接器 可借助 Airbyte 生态连接 300 种数据源完整的可用连接器清单见 连接器总览页。深入看 CSV 连接器 python/pathway/io/csv/init.py 的read签名可以了解文档没有展开的实用参数modestreaming默认或static。streaming 模式下引擎持续监听目录中文件的增、删、改删除文件会把对应行从表中移除static 模式则一次性摄取现有数据两种模式的完整差异见 流式与静态模式文档schema结果表的 Schema也可传None让引擎自动推断csv_settingsCSV 解析器设置分隔符等autocommit_duration_ms默认 1500 毫秒控制流式模式下自动提交更新的节奏是延迟与吞吐之间的直接调优旋钮object_pattern目录内文件过滤模式with_metadata为每行附加_metadataJSON 列包含文件的created_at、modified_at、seen_at等时间戳。5. 转换Transformations用 Python 描述增量计算数据接入后就可以用 Pathway 的转换操作定义处理逻辑。官方示例演示了过滤 分组求和filtered_table input_table.filter(input_table.colA 0) result_table ( filtered_table .groupby(filtered_table.colB) .reduce(sum_valpw.Reducers.sum(pw.this.colC)) )这段代码先保留colA 0的行再按colB分组对每组的colC求和。值得注意的是语义在流式场景下这不是一次性计算而是增量维护的聚合——后续任何新到达或更新的数据到达时Rust 引擎只重算受影响的分组。框架支持的操作面较广概览文档归纳如下类别操作示例算术运算、-、*、/、//、%、**t.select(new_col t.colA t.colB)比较运算、!、、、、t.select(new_col t.colA t.colB)布尔运算AND、\|OR、~NOT、^XORt.select(new_col t.colA (t.colB 3))过滤filtert.filter(pw.this.column value)应用函数pw.applyt.select(new_colpw.apply(func, pw.this.colA))SQL 命令pw.sqlpw.sql(query, tabt)其中pw.apply允许嵌入任意 Python 函数即 UDFpw.sql则提供 SQL 接口见 SQL 文档。基本操作全集见 表操作指南。框架还提供更高级的转换原语Group-by 与聚合见 Group-by / Reduce 手册内置pw.Reducers提供 sum、min、max、mean 等聚合器Join 操作见 Join 手册支持 inner/left/right/outer 等连接模式python/pathway/__init__.py中导出的join、join_inner、join_left、join_right、join_outer即为对应入口。6. 时态转换窗口、ASOF Join 与 Interval Join作为数据流处理框架Pathway 内置一组时态temporal操作概览文档列出了四类类别操作示例窗口操作windowbysliding / tumbling / session 窗口t.windowby(t.time, windowpw.temporal.tumbling(duration...), ...).reduce(...)ASOF now joinasof_now_joint1.asof_now_join(t2, t1.t, t2.t, t1.name t2.name, how..., direction...).select(...)Interval joininterval_joinouter/left/rightt1.interval_join(t2, t1.t, t2.t, pw.temporal.interval(...), t1.col t2.col).select(...)Window joinwindow_joinouter/left/rightt1.window_join(t2, t1.t, t2.t, pw.temporal.sliding(...), t1.col t2.col).select(...)这些算子在源码中位于 python/pathway/stdlib/temporal/ 目录是纯 Python 构建在核心算子之上的标准库实现可以精确定位到每个函数windowby、sliding、tumbling窗口定义python/pathway/stdlib/temporal/_window.pyinterval_join及其 inner/left/right/outer 变体python/pathway/stdlib/temporal/_interval_join.pyasof_now_join及其变体python/pathway/stdlib/temporal/_asof_now_join.py。时态操作的**行为behavior**决定了准确性、延迟、内存消耗三者之间的权衡可以用asof事件到达时才处理与eager数据一到达就提前计算两种策略及其组合来调优详见 时态行为文档 与 窗口行为配置窗口操作本身参见 窗口手册、Interval Join 与 Window Join。ASOF now join 在 索引机制文档 中也有讲解。7. 配置输出把结果送回外部系统管道处理完成后用输出连接器把结果写出框架。最简单的是写 CSV 文件pw.io.csv.write(result_table, ./output/)概览文档列出的常用输出连接器输出连接器示例CSVpw.io.csv.write(table, ./output/)Kafkapw.io.kafka.write(table, rdkafka_settings, topic_nameexample, formatjson)PostgreSQLpw.io.postgres.write(table, output_postgres_settings, sum_table)Google PubSubpw.io.pubsub.write(table, publisher, project_id, topic_id)与输入侧一样输出实现同样对应 python/pathway/io/ 下的各子包kafka、postgres、pubsub、elasticsearch、clickhouse、slack等。从源码结构看每个连接器包通常同时包含read与write两侧能力以支持数据库类系统的读写双向同步底层共用internals/datasource.py与internals/datasink.py中的源/汇抽象。完整清单仍建议查阅 连接器总览页。8. 运行管道pw.run() 与永久运行语义一切就绪后一行命令启动计算pw.run()这里有一个必须建立的认知pw.run()启动的进程会一直监听数据源的新更新直到进程被终止——计算永远运行下去才是框架的正常行为而不是处理完初始数据就退出这一设计动机在 Core Concepts 的用 Rust 引擎运行计算一节有解释。从源码看run函数由 python/pathway/internals/api.py 提供并经 python/pathway/init.py 导出其职责是把前面声明式的管道连接器 转换构成的图交给 Rust 引擎执行。这也解释了 Pathway 的核心架构特征管道定义与计算执行严格分离——你在代码里声明的是处理配方数据源、转换、输出位置真正的数据流动发生在pw.run()之后由 Rust 引擎按增量语义持续推进。9. 延伸方向LLM 工具包与完整学习路径LLM xpackPathway 提供 LLM 扩展包源码位于 python/pathway/xpacks/可在框架内构建嵌入、检索、问答等 LLM 管道详见 LLM xpack 概览入门实操从 第一个实时应用 动手写端到端示例深入概念Core Concepts 完整覆盖连接器、表、转换、输出、数据流与运行时六大核心概念框架对比为什么选择 Live Data Framework 与 批处理模式 可帮助理解它与批式 ETL 的差异。小结一条管道的完整骨架把本文各节拼起来一个最小但完整的 Pathway 实时应用就是五步import pathway as pw # 1. 定义 Schema class InputSchema(pw.Schema): colA: int colB: float colC: str # 2. 输入连接器流式监听 ./data/ 目录 input_table pw.io.csv.read(./data/, schemaInputSchema) # 3. 转换过滤 分组聚合 filtered_table input_table.filter(input_table.colA 0) result_table ( filtered_table .groupby(filtered_table.colB) .reduce(sum_valpw.Reducers.sum(pw.this.colC)) ) # 4. 输出连接器 pw.io.csv.write(result_table, ./output/) # 5. 运行持续监听数据源直到进程终止 pw.run()这套Schema 声明 → 连接器建表 → 转换描述计算 → 输出连接器 →pw.run()的固定骨架正是 Pathway 文档 Live Data Framework 概览 所传递的核心编程范式用静态 Python 代码声明一条会随数据流动而持续更新的实时管道。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表