,以及如何正确等待图任务完成)
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 将 oneTBBoneAPI 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.hcmdline.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()的描述是Callswait_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必须串行并发度 1input_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_waitnode.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_nodeint start_node(g); flow::function_nodeint, int f1(g, flow::unlimited, f1_body{}); flow::function_nodeint, int f2(g, flow::unlimited, f2_body{}); flow::function_nodeint, int f3(g, flow::unlimited, f3_body{}); flow::join_nodestd::tupleint, int join(g); flow::function_nodestd::tupleint, 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_port0(join)); flow::make_edge(f3, flow::input_port1(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),仅供参考