mold 内置 oneTBB Flow Graph 实战:为什么每次都必须调用 wait_for_all(),以及如何正确等待图任务完成
【免费下载链接】moldmold: A Modern Linker 🦠项目地址: https://gitcode.com/GitHub_Trending/mo/mold
导读
本文以 mold 仓库自带的 oneTBB 用户指南文档 always_use_wait_for_all.rst 为核心,系统讲解 Flow Graph(流图)编程中最常见、也最容易引发隐蔽崩溃的一个错误:忘记调用graph::wait_for_all()。读完本文,你将掌握wait_for_all()的完整语义与底层实现原理、为什么销毁图之前必须等待、如何通过try_put_and_wait降低等待延迟,以及后台线程中安全销毁流图的正确姿势,并能直接套用文中给出的可运行代码。
背景:TBB 在 mold 项目中扮演的角色
mold 是一个现代链接器,其核心优化思路之一就是充分利用多核并行。mold 将 oneTBB(oneAPI Threading Building Blocks)作为第三方依赖整体内置在 third-party/tbb 目录中,并在源码中大量使用 TBB 的并行算法:例如 gc-sections.cc 中通过tbb::parallel_for_each并行遍历所有目标文件,arch-arm32.cc 引入<tbb/parallel_for.h>与<tbb/parallel_for_each.h>,cmdline.cc 使用tbb::global_control控制线程数。
在 TBB 提供的众多能力中,Flow Graph 是一套独立的图编程模型:节点(node)是计算单元,边(edge)是消息通道,消息沿边流动并异步调度任务执行。它与parallel_for等普通并行算法最大的不同在于:任务的执行是异步的、由运行时调度器决定的,调用线程并不自动等待任务完成。这就引出了本文的主角wait_for_all()。
核心问题:忘记 wait_for_all() 的典型错误
Flow Graph 编程中最常见的错误之一,就是忘记调用wait_for_all()。always_use_wait_for_all.rst 明确指出:
graph::wait_for_all会阻塞,直到图中派生的所有任务全部完成。这不仅在你需要等待计算结束时有用,而且在销毁 graph 或其任何一个节点之前,调用它都是必须的(necessary)。
文档给出了一个会导致程序失败的经典反例:
void no_wait_for_all() { graph g; function_node< int, int > f( g, 1, []( int i ) -> int { return spin_for(i); } ); f.try_put(1); // program will fail when f and g are destroyed at the // end of the scope, since the body of f is not complete }这段代码在函数作用域结束时,图g与节点f会被析构。然而,用于执行f的 body 的任务(task)此刻可能仍在飞行中(in flight)——try_put(1)只是把消息放入图,任务是否已被调度器执行、是否执行完毕,调用方完全无法确定。
为什么会导致程序失败:悬垂节点与未完成的任务
继续引用文档的解释:当这个未完成的任务最终执行完毕时,它会尝试查找与其节点相连的后继节点(successors),但此时图和节点都已经被从它脚下删除(deleted out from underneath it)。也就是说,任务回调会访问已经被析构的节点与图对象,形成经典的悬垂指针(dangling pointer)访问,后果是未定义行为,轻则数据竞争、重则段错误或随机崩溃。
从源码结构看,这一危险性源于graph析构函数的行为。在 graph_cls.rst 的 API 规范中,~graph()的描述是:
Calls
wait_for_all()on the graph, then destroys the graph.
即析构函数内部会尝试调用一次wait_for_all()再销毁。但问题在于:图中节点同样持有对图上下文的引用,任务也可能引用节点对象本身。如果节点f先被析构(局部变量的析构顺序与构造顺序相反),而图的任务还在执行并要回写或查找该节点,就会访问到已释放的内存。因此文档强调"销毁 graph或任何一个节点之前"都必须手动等待,不能指望析构函数兜底。
再对照 _flow_graph_impl.h 中的注释:
//! Destroys the graph. /** Calls wait_for_all, then destroys the root task and context. */ ~graph();这与规范文档的描述完全一致,可以推断:等待与销毁的时序保证依赖 wait_for_all 被正确调用,而节点先于图析构的顺序使得漏调 wait_for_all 必然引入悬垂访问。
graph::wait_for_all 的语义与底层实现
API 语义
根据 graph_cls.rst 的正式定义:
void wait_for_all();Blocks execution until all tasks associated with the graph have completed or cancelled.
即:阻塞当前线程,直到与该图关联的所有任务完成或被取消。注意这里"完成或被取消"是一个关键语义——即使图被cancel(),wait_for_all也能正常返回,不会永久阻塞。
同一个文档还给出了一系列与wait_for_all相关的状态查询接口,可用于等待结束后的结果判断:
| 接口 | 语义 |
|---|---|
void reset(reset_flags f = rf_reset_protocol) | 按标志重置图(可位或组合,线程不安全,勿并发调用) |
void cancel() | 取消图中所有任务 |
bool is_cancelled() | 最近一次wait_for_all()期间图是否被取消 |
bool exception_thrown() | 最近一次wait_for_all()期间是否有异常抛出 |
底层实现:等待线程会主动"偷活"
wait_for_all()并不是简单地忙等或睡眠。在 _flow_graph_impl.h 中可以读到它的真实实现逻辑:
void wait_for_all() { cancelled = false; caught_exception = false; try_call([this] { my_task_arena->execute([this] { d1::wait(my_wait_context_vertex.get_context(), *my_context); }); cancelled = my_context->is_group_execution_cancelled(); }).on_exception([this] { my_context->reset(); caught_exception = true; cancelled = true; }); ... }关键点有两处:
- 等待被放进
my_task_arena->execute(...)中执行,注释明确写着 "The waiting thread will go off and steal work while it is blocked in the wait_for_all"——阻塞期间,等待线程不会闲下来,而是会去窃取(steal)并执行其他任务。这是 TBB 工作窃取调度器的典型行为:wait_for_all的等待本身就是参与并行计算的过程,因此它不会浪费一个线程。 - 等待结束后,通过
my_context->is_group_execution_cancelled()判断图是否被取消,并把结果记录到cancelled成员;若等待过程中抛出异常,则重置上下文并记录caught_exception = true。这两个成员正是is_cancelled()与exception_thrown()的数据来源。
这从实现层面印证了:wait_for_all()的等待是"活跃等待",也是图生命周期管理不可替代的一环。
正确用法:完整的数据流图示例
Data_Flow_Graph.rst 给出了一个完整且正确的实现:一个生成 1~10 的input_node,两个分别求平方与立方的function_node,以及一个累加求和的function_node(并发度限制为 1,因为它修改共享的sum)。在src.activate()启动数据源之后、读取sum之前,显式调用g.wait_for_all():
class src_body { const int my_limit; int my_next_value; public: src_body(int l) : my_limit(l), my_next_value(1) {} int operator()( oneapi::tbb::flow_control& fc ) { if ( my_next_value <= my_limit ) { return my_next_value++; } else { fc.stop(); return int(); } } }; int main() { int sum = 0; graph g; function_node< int, int > squarer( g, unlimited, [](const int &v) { return v*v; } ); function_node< int, int > cuber( g, unlimited, [](const int &v) { return v*v*v; } ); function_node< int, int > summer( g, 1, & -> int { return sum += v; } ); make_edge( squarer, summer ); make_edge( cuber, summer ); input_node< int > src( g, src_body(10) ); make_edge( src, squarer ); make_edge( src, cuber ); src.activate(); g.wait_for_all(); cout << "Sum is " << sum << "\n"; }注意其中的取舍逻辑,这也是 Flow Graph 并发度设计的直观教材:
squarer与cuber无副作用,可以unlimited并发;summer通过引用修改共享变量sum,必须串行(并发度 1);input_node的 body 通过flow_control& fc的fc.stop()来终止消息源。
没有最后的g.wait_for_all(),cout << "Sum is "完全可能在所有任务完成前就执行,sum的值将不可预测。这正是"计算完成后想读结果就必须等待"的典型场景。
进阶优化:用 try_put_and_wait 等待单条消息
wait_for_all()会等待整个图的所有任务,包括与其他消息无关的工作。在低延迟敏感、逐消息处理的场景下,这可能引入不必要的等待。try_put_and_wait.rst 介绍了一个预览特性接口try_put_and_wait:
node.try_put_and_wait(msg)在节点上执行node.try_put(msg),并等待与该msg相关的所有工作完成。相比graph::wait_for_all,它可以降低延迟,因为后者会等待包括与输入消息无关在内的所有工作。
启用方式(二选一)并包含头文件:
#define TBB_PREVIEW_FLOW_GRAPH_FEATURES // macro option 1 #define TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT // macro option 2 #include <oneapi/tbb/flow_graph.h>该接口保证两条语义:
- 图中任意节点因处理
msg(或由其派生的中间结果)而创建的任务全部完成; - 由
msg派生的中间结果不再残留在图中任何缓冲区中。
仓库自带的完整示例见 try_put_and_wait_example.cpp,其核心结构是broadcast_node分发给两条并行链(f1→f2与f3),经join_node汇合后由f4消费,然后在parallel_for内逐条提交并立即等待:
flow::graph g; flow::broadcast_node<int> start_node(g); flow::function_node<int, int> f1(g, flow::unlimited, f1_body{}); flow::function_node<int, int> f2(g, flow::unlimited, f2_body{}); flow::function_node<int, int> f3(g, flow::unlimited, f3_body{}); flow::join_node<std::tuple<int, int>> join(g); flow::function_node<std::tuple<int, int>, int> f4(g, flow::serial, f4_body{}); flow::make_edge(start_node, f1); flow::make_edge(f1, f2); flow::make_edge(start_node, f3); flow::make_edge(f2, flow::input_port<0>(join)); flow::make_edge(f3, flow::input_port<1>(join)); flow::make_edge(join, f4); // Submit work into the graph parallel_for(0, 100, & { start_node.try_put_and_wait(input); // Post processing the result of input });注意该特性的两个重要限制(文档中以 caution 标注):
- 不要在流图末端使用缓冲类节点(如
buffer_node、queue_node、overwrite_node等):最终结果不会被自动消费,try_put_and_wait可能无限等待;对overwrite_node需要显式调用clear()或覆盖新值,对write_once_node需要显式clear()。 multifunction_node与async_node暂不支持:图中包含它们时,try_put_and_wait可能在初始消息的计算仍在进行时就提前返回。
此外,该接口并不是所有节点的专利——continue_node、function_node、overwrite_node、write_once_node、buffer_node、queue_node、priority_queue_node、sequencer_node、limiter_node、broadcast_node、split_node均提供同名重载,完整签名与各自 Effects 见 try_put_and_wait.rst。
相关陷阱:后台线程中的图销毁
有时你不想让主线程被wait_for_all()阻塞,但销毁图之前调用 wait_for_all 仍然是最安全的做法。TBB 用户指南在 destroy_graphs_outside_main_thread.rst 中给出的推荐方案是:把"建图 → 喂数据 → 等待"整体封装成一个任务,通过task_arena::enqueue投递到后台执行:
class background_task { public: void operator()() { graph g; function_node< int, int > f( g, 1, []( int i ) -> int { return spin_for(i); } ); f.try_put(1); g.wait_for_all(); } }; void no_wait_for_all_enqueue() { task_arena a; a.enqueue(background_task()); // do other things without waiting… }文档同时提醒:被投递的任务何时执行是不确定的,如果程序需要用到该任务的结果,甚至要确保它在程序结束前完成,就必须通过某种信号机制(如std::promise、原子标志或a.wait())与主线程同步。
这三条"等待与销毁"建议在用户指南中被组织为一个独立的主题章节,见 Flow-Graph-waiting-tips.rst,它与avoid_dynamic_node_removal、destroy_graphs_outside_main_thread共同构成 Flow Graph 生命周期管理的完整提醒。
排查建议:神秘行为的第一个检查点
原文档在结尾给出了一条极其实用的排障建议:
如果你使用了 Flow Graph 并看到了难以解释的行为,首先检查你是否调用了
wait_for_all()。
这条建议可以展开为一份简要的排查清单:
- 图中任务是否确实全部完成?未完成的节点 body 可能写入了未初始化的数据。
- 图或节点是否在任务完成前被析构?这是崩溃与随机错误的头号来源。
- 数据是否还残留在节点缓冲区中?缓冲类节点(
buffer_node等)不会自动消费数据。 - 是否误用了
try_put_and_wait而不满足其前置条件(末端缓冲节点、multifunction_node/async_node)? - 是否在非主线程销毁图而没有先等待?
小结
graph::wait_for_all()是 Flow Graph 编程中生命周期管理的第一道防线:它阻塞至图中所有任务完成或取消,等待期间线程会主动窃取工作以继续参与计算(见 _flow_graph_impl.h),并且在图与节点析构之前调用是强制要求而非可选项。对于逐消息的低延迟场景,可以改用预览接口try_put_and_wait将等待范围收窄到单条消息;对于后台线程,则应将"建图—等待—销毁"整体封装进被投递的任务。记住文档的核心结论:用 Flow Graph 遇到诡异行为,先检查wait_for_all调用了没有。
【免费下载链接】moldmold: A Modern Linker 🦠项目地址: https://gitcode.com/GitHub_Trending/mo/mold
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考