news 2026/9/8 23:47:09

Timely Dataflow 入门实战:从 hello 示例看懂流式计算中的 Worker、Exchange 与并行执行

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Timely Dataflow 入门实战:从 hello 示例看懂流式计算中的 Worker、Exchange 与并行执行

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) }

它的工作流程可以概括为三步:

  1. Config::from_args解析命令行参数,得到通信配置;
  2. 调用execute(config, func)依据配置初始化通信并创建 worker 线程;
  3. 在每个 worker 上运行你传入的闭包func(&mut Worker<...>),最终返回WorkerGuards<T>(可用于 join 收集各 worker 的返回值或错误)。

闭包签名中的约束T: Send + 'staticF: 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(); } }

每一步做三件事:

  1. input.send(round):仅当index == 0时向输入发送整数round
  2. input.advance_to(round + 1):把输入的时间戳推进round + 1,表明这个时间之前的输入已经全部结束;
  3. while probe.less_than(input.time()) { worker.step(); }:只要探测点的进度还落后于当前输入时间(即还有未处理完的数据),就反复调用worker.step()让运行时继续执行算子。

这正是 timely dataflow"时间驱动"模型的缩影:数据不是以"发送完立即处理完"的方式流动,而是被贴上逻辑时间戳,系统通过advance_to告知"该时间戳之前不会再有新数据",各算子依据进度信息决定何时可以安全地完成某一时刻的处理。

数据混洗的规则

-w 2这类多 worker 场景下,worker 0 注入全部0..10exchange(|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 == 1threads > 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 -p1

shell 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

两个要点:

  1. 每个进程只打印自己处理的那一半——shell B 里看不到偶数,shell A 里看不到奇数,数据被按值划分到两个进程中;
  2. 首个进程会等待同伴就绪——这是分布式启动的握手阶段,-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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/8 23:46:05

快手接口签名sig、sig3与NStoken原理及测试用例详解

简介&#xff1a;快手sig3、sig、NStoken算法资源面向移动应用开发者、爬虫及逆向分析人员&#xff0c;聚焦快手接口交互中的请求签名与身份令牌生成问题。压缩包共4个文件&#xff0c;包含2个Python脚本和2个数据文件&#xff08;data0、data1&#xff09;&#xff0c;脚本分别…

作者头像 李华
网站建设 2026/9/8 23:43:15

Chiplet异构集成封装设计:芯片-封装-PCB协同与信号完整性实践

最近不少做封装设计的朋友问我同一个问题&#xff1a;小芯片集成&#xff08;Chiplet&#xff09;明明是芯片设计的事&#xff0c;为什么我们这些做 PCB、做封装的人反而比芯片团队还忙&#xff1f;这其实问到了点子上。异构集成电路封装设计里的 Chiplet 集成&#xff0c;核心…

作者头像 李华
网站建设 2026/9/8 23:41:40

Goose 本地部署实操指南:3 步装好 CLI 并跑通第一个智能体任务

Goose 本地部署实操指南&#xff1a;3 步装好 CLI 并跑通第一个智能体任务 【免费下载链接】goose an open source, extensible AI agent that goes beyond code suggestions - install, execute, edit, and test with any LLM 项目地址: https://gitcode.com/GitHub_Trendin…

作者头像 李华
网站建设 2026/9/8 23:31:27

STM32H743高性能MCU实战:从选型到量产的全流程解析

最近帮客户做一套工业视觉检测的预处理板&#xff0c;主控选型的时候纠结了很久。一开始想用MPU加Linux的方案&#xff0c;但考虑到成本、功耗和现场环境&#xff0c;最后还是回到了高端MCU这条路上。在对比了NXP的RT1170、Microchip的SAMA7G54和ST的STM32H743之后&#xff0c;…

作者头像 李华