Timely Dataflow 入门实战:从 hello 示例看懂流式计算中的 Worker、Exchange 与并行执行
【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway
本文以 Timely Dataflow 官方指南 A Simple Example 为主体骨架,以本仓库(Pathway 项目,其底层基于 timely-dataflow / differential-dataflow 构建)内随附的 timely-dataflow 源码与示例为辅证,逐行拆解第一个完整的 dataflow 程序hello,并演示单线程、单进程多线程、多进程三种运行方式,帮助你理解 worker 模型、exchange 数据混洗、时间戳推进与 probe 探测等核心机制,为后续阅读更复杂的数据流系统(如 Pathway 本身)打下基础。
从一个 "hello" 程序说起
Timely Dataflow 的目标是捕捉大量数据处理惯用法(idiom),因此很难用一个例子同时展示全部特性。官方指南刻意挑选了一个只覆盖核心功能的最小示例:程序初始化一个 timely dataflow 计算,参与者可以向其中喂入一串数字,这些数字会根据其取值在 worker 之间交换(exchange),每个 worker 在"看到"数字时打印到屏幕。
本仓库中随附的源码副本位于 examples/hello.rs,完整的可运行程序如下:
extern crate timely; use timely::dataflow::InputHandle; use timely::dataflow::operators::{Input, Exchange, Inspect, Probe}; fn main() { // initializes and runs a timely dataflow. timely::execute_from_args(std::env::args(), |worker| { let index = worker.index(); let mut input = InputHandle::new(); // create a new input, exchange data, and inspect its output let probe = worker.dataflow(|scope| scope.input_from(&mut input) .exchange(|x| *x) .inspect(move |x| println!("worker {}:\thello {}", index, x)) .probe() ); // introduce data and watch! for round in 0..10 { if index == 0 { input.send(round); } input.advance_to(round + 1); while probe.less_than(input.time()) { worker.step(); } } }).unwrap(); }这段代码虽然不到四十行,却涵盖了 timely dataflow 应用中最核心的五个概念:入口执行函数、worker 与输入句柄、dataflow 描述、操作符链、进度与探测。下面逐层展开。
入口:execute_from_args
程序从timely::execute_from_args(std::env::args(), |worker| ...)开始。该函数的真实实现位于 execute.rs:
#[cfg(feature = "getopts")] pub fn execute_from_args<I, T, F>(iter: I, func: F) -> Result<WorkerGuards<T>,String> where I: Iterator<Item=String>, T:Send+'static, F: Fn(&mut Worker<Allocator>)->T+Send+Sync+'static, { let config = Config::from_args(iter)?; execute(config, func) }它的工作流程可以概括为三步:
- 用
Config::from_args解析命令行参数,得到通信配置; - 调用
execute(config, func)依据配置初始化通信并创建 worker 线程; - 在每个 worker 上运行你传入的闭包
func(&mut Worker<...>),最终返回WorkerGuards<T>(可用于 join 收集各 worker 的返回值或错误)。
闭包签名中的约束T: Send + 'static、F: Send + Sync + 'static意味着:你的业务闭包与返回值都必须是可在线程/进程间传递的,这是分布式安全性的第一道编译期保障。
Worker 与InputHandle
闭包收到的是一个&mut Worker,每个 worker 拥有独立的index()(从 0 开始编号)。代码用worker.index()取出本 worker 的编号,随后创建InputHandle::new()作为输入句柄。
在单进程多线程模式下,一次execute_from_args会创建-w个 worker 线程;每个线程的闭包都会被执行一次,因此代码中必须通过if index == 0之类的判断来决定谁负责注入数据,否则每个 worker 都会各自发送0..10,数据总量就会翻倍。
描述 dataflow:worker.dataflow
let probe = worker.dataflow(|scope| scope.input_from(&mut input) .exchange(|x| *x) .inspect(move |x| println!("worker {}:\thello {}", index, x)) .probe() );worker.dataflow(...)用于在 worker 内部描述(声明式构建)一个数据流计算图。图中依次连接四个操作符:
input_from(&mut input):把外部的InputHandle暴露为数据流的输入源。注意input是在 dataflow 外部创建的,通过引用注入,这样外层循环才能持续向数据流喂数据;exchange(|x| *x):按闭包返回的键值把记录重新分发。这里*x即数字本身,因此所有相同的数字会被路由到同一个 worker;inspect(move |x| ...):对每条数据执行副作用(打印),闭包用move捕获index,把"我在哪个 worker"的信息移入算子;probe():在流末尾放置一个探测点,返回一个ProbeHandle,供外层判断进度。
外层驱动循环:数据、时间戳与进度
真正"推动"计算的是main中 dataflow 之外的循环:
for round in 0..10 { if index == 0 { input.send(round); } input.advance_to(round + 1); while probe.less_than(input.time()) { worker.step(); } }每一步做三件事:
input.send(round):仅当index == 0时向输入发送整数round;input.advance_to(round + 1):把输入的时间戳推进到round + 1,表明这个时间之前的输入已经全部结束;while probe.less_than(input.time()) { worker.step(); }:只要探测点的进度还落后于当前输入时间(即还有未处理完的数据),就反复调用worker.step()让运行时继续执行算子。
这正是 timely dataflow"时间驱动"模型的缩影:数据不是以"发送完立即处理完"的方式流动,而是被贴上逻辑时间戳,系统通过advance_to告知"该时间戳之前不会再有新数据",各算子依据进度信息决定何时可以安全地完成某一时刻的处理。
数据混洗的规则
在-w 2这类多 worker 场景下,worker 0 注入全部0..10,exchange(|x| *x)会把它们按值打散。文档明确指出:唯一保证是"exchange 闭包求值结果相同的记录一定到达同一个 worker";当前实现实际采用"数字对 worker 数取余"的路由方式。所以偶数归 worker 0、奇数归 worker 1,输出出现worker 0: hello 0 / worker 1: hello 1 / worker 0: hello 2 ...交替的形态。这种按取余的简单路由正是分布式流式系统做数据分区(partition)的常见雏形——Pathway 与 differential-dataflow 中的数据分组/关联同样依赖这类 key 路由语义。
单线程运行:默认配置的零成本启动
先克隆并编译(文档中的命令面向独立仓库;在本仓库中该 crate 已内嵌于 external/timely-dataflow 目录,对应源码在 timely/ 与 communication/ 下):
cargo build cargo build --example hello cargo run --example hello第三个命令即可一步完成构建并运行(Rust 会自动完成必要的编译),输出为:
worker 0: hello 0 worker 0: hello 1 worker 0: hello 2 worker 0: hello 3 worker 0: hello 4 worker 0: hello 5 worker 0: hello 6 worker 0: hello 7 worker 0: hello 8 worker 0: hello 9这是最基本的配置:一个进程内只有一个 worker 线程。从配置解析源码 communication/src/initialize.rs 可以看到,当既未指定多进程也未指定多线程时,Config::from_matches会返回Config::Thread这一最简单的通信变体,try_build时直接使用进程内线程级分配器,连跨线程通道都省掉了。换言之,单 worker 是 timely 刻意保留的"零开销基线"。
单进程多线程:-w / --workers参数
用-w(或长选项--workers)指定进程内 worker 线程数。注意cargo run与程序参数之间必须加--,否则-w2会被 cargo 误读:
cargo run --example hello -- -w2输出示例:
worker 0: hello 0 worker 1: hello 1 worker 0: hello 2 worker 1: hello 3 worker 0: hello 4 worker 1: hello 5 worker 0: hello 6 worker 1: hello 7 worker 0: hello 8 worker 1: hello 9对比单线程输出,唯一的可见差异是worker index 交替出现——这正是"存在多个 worker 并且处理了不同数据"的证据。worker 0 注入全部数据(由代码中的index == 0守卫保证,否则每个 worker 都会发一遍0..10),然后数字经 exchange 在 worker 间被重新分布。
配置层的对应逻辑同样在Config::from_matches:当processes == 1且threads > 1时返回Config::Process(threads);若同时指定-z/--zerocopy则升级为Config::ProcessBinary,为进程内通信启用零拷贝优化。其余两个选项-r/--report会打印连接/进度报告,可在多 worker 场景下观察各线程的同步情况。
多进程:-n、-p与-h
要让计算跨越多个进程(甚至多台机器),需要-n与-p两个参数协同:
-n / --processes:声明本次计算总共有多少个进程参与;-p / --process:声明当前进程是其中的第几个(从 0 开始)。
参数解析位于 communication/src/initialize.rs 的install_options:
opts.optopt("w", "threads", "number of per-process worker threads", "NUM"); opts.optopt("p", "process", "identity of this process", "IDX"); opts.optopt("n", "processes", "number of processes", "NUM"); opts.optopt("h", "hostfile", "text file whose lines are process addresses", "FILE"); opts.optflag("r", "report", "reports connection progress"); opts.optflag("z", "zerocopy", "enable zero-copy for intra-process communication");多进程模式同样支持-w指定每个进程内的线程数,因此实际可以形成"多个进程 × 每进程多线程"的完整并行矩阵。当processes > 1时:
- 若提供了
-h <hostfile>,从文件中逐行读取前processes个地址(形如host:port)作为进程地址表;若读取到的地址数不足processes,会返回错误; - 若省略
-h,默认回退到本机地址localhost:2101 + index,即默认在本机协调多个进程; - 每个进程的 worker 数由
-w决定,进程身份由-p决定。
该模式最终构造Config::Cluster,并通过网络初始化真正的跨进程通信通道。
开两个 shell 跑一个分布式程序
在两个 shell 中分别启动同一命令、仅改变-p:
# shell A(会先阻塞,等待其他进程就绪) cargo run --example hello -- -n2 -p0 # shell B cargo run --example hello -- -n2 -p1shell B 的输出:
worker 1: hello 1 worker 1: hello 3 worker 1: hello 5 worker 1: hello 7 worker 1: hello 9回到 shell A,可看到进程已经运行并输出了另一半结果:
worker 0: hello 0 worker 0: hello 2 worker 0: hello 4 worker 0: hello 6 worker 0: hello 8两个要点:
- 每个进程只打印自己处理的那一半——shell B 里看不到偶数,shell A 里看不到奇数,数据被按值划分到两个进程中;
- 首个进程会等待同伴就绪——这是分布式启动的握手阶段,
-n2 -p0的 shell A 会先"挂起",直到第二个进程接入。文档中对默认 host 的说明与源码中addresses.push(format!("localhost:{}", 2101 + index))的逻辑完全一致。
若想跨机器运行,把-h指向一个 host 文件即可。文档中 host 文件格式为每行一个地址:端口,与本仓库install_options中"text file whose lines are process addresses"的说明吻合:
host0:port host1:port host2:port host3:port对应每个主机分别执行cargo run -- -w 2 -n 4 -h hosts.txt -p 0/1/2/3,即可得到"4 进程、每进程 2 worker、共 8 worker"的分布式计算。可见execute_from_args之所以接受std::env::args(),正是为了让同一套程序能透明地适配单线程、多线程、多进程全部三种部署形态。
本示例在更大图景中的位置
hello虽然简单,却已经铺好了 timely dataflow 后续章节的全部伏笔:chapter_1将展开时间戳(timestamps)与进度(progress)的形式化定义;chapter_2讲如何创建输入、观察输出与手写自定义算子;chapter_3讲输入驱动与 probe 监控;chapter_4则引入 scope、迭代与背压等高级话题。本书章节目录可见 SUMMARY.md。
对当前仓库的读者来说,这个例子的价值还在于:Pathway 正是构建在 Rust 流式计算基石之上的 Python ETL / 实时分析 / RAG 框架,其依赖的及时数据流与差分数据流能力与 external/timely-dataflow、external/differential-dataflow 中随附的源码同源。理解本页的 worker 模型、exchange 路由、advance_to时间推进与 probe 进度探测,等于掌握了阅读 Pathway 底层执行引擎的一把钥匙——当你看到 Pathway 中"数据在多 worker 间 shuffle、按时间戳批处理、输出可探测"等行为时,本质上都在与本文演示的这套机制打交道。
【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考