深入解析 Pathway 多 Worker 架构:分片执行、进度同步与分布式部署

深入解析 Pathway 多 Worker 架构:分片执行、进度同步与分布式部署 深入解析 Pathway 多 Worker 架构分片执行、进度同步与分布式部署【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway本文基于 docs/2.developers/4.user-guide/80.advanced/10.worker-architecture.md 展开并结合仓库内 CLI 源码、引擎实现与相邻文档进行佐证扩充。Pathway 是一套 Python ETL 框架面向流式处理、实时分析、LLM 数据管道与 RAG 场景。其底层是 Rust 编写的流式数据流引擎而多 Workerworker并行执行正是这套架构能够把同一个 Python 程序横向扩展到多线程、多进程乃至多台机器上的关键机制。读完本文你将理解 Worker 如何通过唯一 ID 划分数据分片、如何在彼此之间交换数据与进度信息以维持结果一致性、本地状态如何快照并用于故障恢复以及如何用pathway spawn和--addresses/--process-id在单机或多机上真正把管道跑起来。为什么需要理解 Worker 架构Pathway 的使用体验是用 Python 描述计算由引擎负责执行。用户写一段声明式的数据流程序交给框架后框架会对计算图进行优化再转换为底层 Rust 引擎可执行的基础算子文档原文。如果你只在一台机器上以单线程方式运行小型管道通常不需要关心 Worker 的存在但一旦遇到以下情况Worker 架构就变得至关重要数据量或中间状态超出单机内存单机 CPU 核数被计算密集型的实时任务打满希望让处理逻辑就近运行在分区的数据源如 Kafka 分区旁边以减少网络传输需要在大流量/低流量波动下自动扩缩容。在上述场景中并行度与一致性都直接由 Worker 的数量、分片方式和通信机制决定。核心设计每个 Worker 都运行同一份数据流Worker 架构的第一条设计原则是每个 Worker 都运行相同的数据流但各自处理不同子集shard/partition的数据。具体表现为每个 Worker无论是一个线程还是一个进程都执行完全相同的 Python 脚本构建出完全相同的底层数据流。这让一套代码、多处执行成为可能——你并不需要为不同机器编写不同的数据读取逻辑。每个 Worker 拥有自己的身份一个唯一的连续编号Worker ID。Worker 用这个 ID 决定哪些数据分片归自己负责。对于支持显式分区的数据源例如分区后的 Kafka TopicWorker 只读取自己所属分区的数据实现天然的并行读取。这一点在引擎侧也有印证Rust 引擎在构建数据流时通过scope.index()获取当前 Worker 在全体 Worker 中的序号见 src/engine/dataflow.rs数据源侧则使用source_reader_worker_id记录负责读取该数据源的 Worker 编号并用位掩码统计哪些读取 Worker 处于空闲/活跃状态见 src/connectors/synchronization.rs这些机制共同支撑了按 Worker 分片读取与同步的上层语义。分片数据源如何被并行读取并非所有数据源都自带分区概念Pathway 对它们做了分级处理数据源类型读取方式说明显式分区源如分区 Kafka Topic按分区分配给 Worker每个 Worker 只读取属于自己 ID 对应责任范围的分区无显式分区但可拆分的源如文件系统连接器稳定哈希按文件分配对每个文件的路径做稳定哈希将其恰好分配给一个 Worker从而让读取量随 Worker 数量线性扩展其余无法拆分的源单 Worker 代读后转发由某一个 Worker 负责读取数据再把数据片段转发给其他 Worker其中文件系统连接器按路径稳定哈希分配文件这一策略尤其值得注意它不需要元数据服务协调只要所有 Worker 对同一路径算出相同的哈希就能保证每个文件恰好被一个 Worker 读到不会出现重复或遗漏同时天然随 Worker 数量扩展读取吞吐。这也是从源码结构可以推断出的设计取向——数据源与 Worker 的绑定关系被集中管理在引擎的同步逻辑中src/connectors/synchronization.rs。Worker 间通信与进度交换多 Worker 并行只是手段让并行结果保持一致才是目的。Pathway 的 Worker 之间会在需要时互相通信通信方式取决于 Worker 的实现形态同一进程内的线程 Worker通过共享内存通信同一台机器上的进程 Worker通过**套接字socket**通信不同机器上的 Worker通过网络上的 TCP 等套接字通信。在数据交换之外更关键的是进度progress信息交换。数据流中的每个节点都会跟踪自己的进度并借助数据流的拓扑结构在处理完一段输入数据后高效地通知其下游/对等节点。这套机制保证了一个重要的一致性性质Pathway 产出的每一个结果都只依赖于输入数据流的一个已知前缀。换句话说某个结果被产出时系统明确知道它对应的输入到达了哪个位置——这正是增量/流式计算得以正确的基础也是后续从快照回放重算能力的根基。文档明确说明Pathway 的基础数据流设计概念沿袭了微软 Naiad 系统的奠基性工作SOSP 2013其通信原语、系统内部的时间概念以及基于内存的状态表示则建立在 external/timely-dataflow 与 external/differential-dataflow 之上——该仓库的external/目录下即保留了这两套 Rust 数据流库的源码副本可对照阅读底层实现。从 CLI 看 Worker 的形态线程与进程要让程序真正以多 Worker 形态运行需要借助pathway spawn命令。从 python/pathway/cli.py 可以看到spawn命令支持的参数定义参数默认值含义-t, --threads N1每个进程内运行的线程数-n, --processes N缺省单进程进程数与--addresses互斥--first-port PORT10000通信使用的起始端口使用--addresses时被忽略--addresses host0:port0,host1:port1无多机部署时各进程的host:port列表进程数由列表长度推导-pi, --process-id N无当前机器上本进程在--addresses列表中的下标配合--addresses使用其中每个进程都可以再包含多个线程因此Worker 数量实际等于进程数 × 每进程线程数。运行方式很简单# 单机、单进程默认 python pipeline.py # 单机、以 2 个进程运行 pathway spawn -n 2 python pipeline.py # 单机、1 个进程内跑 4 个线程 pathway spawn -t 4 python pipeline.py # 单机、2 个进程、每进程 2 个线程共 4 个 Worker pathway spawn -n 2 -t 2 python pipeline.py单机多进程模式与多机模式的一个关键差异是单机模式下pathway spawn会自行在127.0.0.1上分配通信端口并负责拉起所有进程而多机模式下需要你为每个机器手动启动一个进程并告诉每个进程其余进程都在哪里。cli.py中的参数校验逻辑validate_and_resolve_spawn_args见 python/pathway/cli.py会强制保证--threads、--processes至少为 1--processes与--addresses互斥——设置了地址列表时进程数由列表长度推导设置了--addresses时--process-id必填且其取值必须在0 .. 地址数-1之间--addresses中不允许出现重复条目端口必须是合法范围。状态存储、快照与故障恢复每个 Worker 的职责不只是执行算子数据流图中有状态算子stateful operators还拥有各自本地的状态存储以及它负责的那部分输入输出连接器。为保证在故障后能够恢复而不从头重算每个 Worker 会异步地把状态保存到一个持久位置例如 S3。当任一环节失败时所有 Worker 共同确定它们各自写入的最后一个快照全体回退rewind到该快照对应的计算位置从该位置继续执行而不是丢弃全部进度。在代码侧持久化通过pw.persistence.Config声明例如使用本地文件系统作为后端from pathway.internals import api import pathway as pw persistence_config pw.persistence.Config( backendpw.persistence.Backend.filesystem(your_persistent_storage_path), # 例如 /tmp/Pathway-Cache persistence_modeapi.PersistenceMode.OPERATOR_PERSISTING, ) pw.run(persistence_configpersistence_config)关于多 Worker 并行、持久化与扩缩容配置的完整实操可参考同目录下的进阶文档 60.worker_count_scaling.md动态扩缩容与worker_scaling_enabled等参数与 30.consistency.md。分布式部署多机、Kubernetes 与固定地址池当单机无法满足要求时可以把管道分布到多台机器。每台机器运行一个进程它们共同构成一个逻辑上完整的计算。不同机器上的 Worker 通过 TCP 通信交换数据与进度信息的方式与同机进程完全一致。文档给出了一条明确的部署约束多机分布式部署可以使用 Kubernetes 及其云实现AKS、EKS并且 Pathway 假定使用StatefulSet有状态副本集部署所有 Pod 都在线是计算成功运行的前提。需要说明的是文档同时指出多机生产级分布式部署属于 Pathway Enterprise 版本能力范畴并支持与现有 Helm Chart 及 k8s 工具链集成非企业版多机运行通常需要获得相应的 Scale/Enterprise 许可。从仓库的相邻文档 70.running_on_multiple_machines.md 可以还原出具体部署步骤例如两台机器192.168.1.10与192.168.1.11的启动命令# 机器 0 pathway spawn \ --addresses 192.168.1.10:9000,192.168.1.11:9000 \ --process-id 0 \ python pipeline.py # 机器 1 pathway spawn \ --addresses 192.168.1.10:9000,192.168.1.11:9000 \ --process-id 1 \ python pipeline.py其中关键规则是所有进程收到的--addresses列表顺序必须完全一致每条命令中只有--process-id不同取值从0覆盖到列表长度 - 1--process-id指定本进程在地址列表中的下标进程 0 绑定192.168.1.10:9000进程 1 绑定192.168.1.11:9000各进程可以任意顺序启动先启动的进程会等待其他进程全部连上后再开始计算此时日志会打印 Preparing Pathway Live Data Framework computation一台机器也可以承载多个进程只需同一 host 使用不同端口如192.168.1.10:9000,192.168.1.10:9001,192.168.1.11:9000--threads与--addresses相互独立可叠加使用例如--threads 2表示每个机器进程内再开两个线程。分布式部署的注意点与限制建议使用共享存储做持久化跨机器运行时若任一进程崩溃整个集群都必须重启。持久化如 S3/GCS/Azure Blob/NFS 这类所有机器可读写的共享后端能保证管道从上一个检查点恢复而不是从头回放。可参考 70.running_on_multiple_machines.md 中给出的 S3 后端配置示例。固定地址池不支持动态扩缩容使用--addresses时进程集合在运行期间是固定的系统无法在运行时增删机器Worker 发出的扩缩容信号会被忽略并产生警告日志。需要动态扩缩容时应改用--processes让 Pathway 自行管理进程相关机制详见 60.worker_count_scaling.md。所有进程必须全部就绪存在一个进程启动失败或过慢时其他进程会无限期等待没有部分启动或降级模式。至少一次at-least-once投递语义崩溃后恢复会从最近已提交的检查点重放数据检查点之后、崩溃之前写入的记录可能被再次处理。恰好一次exactly-once语义属于企业版能力。全集群须版本一致所有机器必须运行相同版本的 Pathway 与相同的管道代码版本不匹配会导致连接失败或未定义行为。网络要求每台机器都必须能通过指定端口访问其他所有机器Pathway 不支持在 Worker 之间做 NAT 穿透或经由代理转发。总结Worker 架构的关键事实清单同码异构分片每个 Worker 运行同一份数据流脚本用唯一 Worker ID 决定自己负责哪些分片源数据并行读取的前提是分片归属明确。分级并行读取显式分区源按分区读文件系统类源按路径稳定哈希单归属分配无法拆分的源由单个 Worker 代读再转发。通信与进度同步线程走共享内存、进程/跨机走套接字节点按拓扑传播处理进度保证每个输出都对应输入流的已知前缀这是结果一致性的根基。本地状态 异步快照有状态算子持有本地状态存储Worker 异步把状态持久化到远端如 S3故障时全体回退到最近快照重放恢复。单机扩展用spawn -n/-t多机部署用--addresses/--process-id分布式以 StatefulSet 方式部署于 KubernetesAKS/EKS企业版提供生产级多机支持。固定地址池限制--addresses模式进程数恒定、不可动态扩缩容动态扩缩容需启用持久化并以--processes模式启动。想要从代码层面进一步验证上述机制可重点阅读 python/pathway/cli.pyspawn 参数解析与多进程拉起、src/connectors/synchronization.rs数据源读取与 Worker 同步以及external/目录下的 timely-dataflow / differential-dataflow 源码底层通信与时间推进原语。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考