多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

oneTBB 流图节点到任务的映射:理解 function_node 如何生成与调度任务

oneTBB 流图节点到任务的映射:理解 function_node 如何生成与调度任务 并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载流图Flow Graph是 oneTBB 中用于表达数据流与依赖关系的核心抽象程序把计算组织成由节点node与边edge构成的图消息沿有向边在节点之间流动节点内部则以 oneTBB 任务task的形式异步执行函数体。本文以官方用户指南中Mapping Nodes to Tasks一节的执行时间线为核心逐帧拆解一个典型的两节点function_node示例在try_put与wait_for_all之间的完整生命周期并结合 flow_graph.h 与 _flow_graph_node_impl.h 中的源码讲清楚节点并发度如何决定任务生成节奏消息如何缓冲wait_for_all为什么不会空转这几个关键问题。读完本文你将能根据节点并发度准确预测流图的行为时序并学会用任务池视角解释流图的执行。从上一节的例子说起两节点图的完整执行在用户指南的 Nodes 与 Edges 两节中我们构造过这样一个流图节点n以unlimited并发度接收三个输入1、2、3对每个输入先打印、再spin_for(v)忙等v秒、再打印并原样返回节点m以并发度 1 接收n的输出平方后同样打印、忙等、再打印。二者之间用make_edge(n, m)建立有向边#include oneapi/tbb/flow_graph.h using namespace oneapi::tbb; graph g; function_node int, int n( g, unlimited, []( int v ) - int { cout v; spin_for( v ); cout v; return v; } ); function_node int, int m( g, 1, []( int v ) - int { v * v; cout v; spin_for( v ); cout v; return v; } ); make_edge( n, m ); n.try_put( 1 ); n.try_put( 2 ); n.try_put( 3 ); g.wait_for_all();正文开头的执行时间线图展示了这段代码的一种可能运行结果系统线程充足时。下文将按时间顺序逐段拆解图中发生的每一类事件。三次 try_put 生成三个任务图中标注的三个n.try_put(1)、n.try_put(2)、n.try_put(3)是图的操作栈一侧发生的事件。因为节点n以unlimited构造每来一条消息运行时都会立即创建一个执行λn的任务任务λn(1)对消息 1 执行函数体任务λn(2)对消息 2 执行函数体任务λn(3)对消息 3 执行函数体。源码层面的对应逻辑在 _flow_graph_node_impl.htry_put_task_impl遇到my_max_concurrency 0即unlimited的取值 0见 flow_graph.h 中enum concurrency { unlimited 0, serial 1 }时直接走create_body_task路径为每条输入分配一个apply_body_task_bypass任务并返回。换句话说unlimited语义就是来一条消息、建一个任务。为什么 λn 的三个任务可以并发执行任务创建出来之后会被投放到 oneTBB 的任务池work pool。三个任务之间没有相互依赖因此只要系统中有足够的线程它们就可以同时执行。这正是unlimited并发度的含义节点不限制自身在飞行in-flight的任务数量。这里需要强调一个容易混淆的概念生成任务 ≠ 创建线程。unlimited只是允许节点无限制地生成任务任务最终仍由 oneTBB 线程池中数量有限的线程来执行。因此图中展示的是有足够线程时的理想时间线——若系统中只有一个线程这三个任务仍会依次排队串行执行。用户指南在 Nodes 一节中同样提醒尽管图可能生成大量任务但只有线程池中可用的那几条线程会真正执行它们。每完成一个任务结果就沿边送入 m当n的某个任务执行完λn得到输出v之后运行时会把该输出通过make_edge(n, m)建立的边自动投递给后继节点m图中标注为m.try_put(1)、m.try_put(2)、m.try_put(3)。源码中这一步发生在apply_body_impl_bypass里apply_body_impl 调用函数体tbb::detail::invoke(*my_body, i)得到输出值vapply_body_impl_bypass 随后调用successors().try_put_task(v)把v投递给后继节点集合function_node的后继以broadcast_cache维护见 flow_graph.h。由于n的三个任务可能并行完成三条消息到达m的顺序在真实运行中并不保证与try_put的调用顺序一致——时间线图只是展示了按到达顺序处理这一种可能。并发完成的顺序取决于任务调度与各次spin_for的时长。并发度 1 的 m串行执行与消息缓冲与n不同节点m以并发度1构造因此它不会为每条到达的消息立即生成任务。相反m会在消息到达的顺序上串行地依次生成任务执行完λm(1)之后才生成执行λm(2)的任务依此类推。这就是时间线图中λm(1) → λm(2) → λm(3)严格串行的原因。源码中的并发度记账机制function_input_base用两个计数器管理并发my_max_concurrency构造时指定的并发上限serial时为 1my_concurrency当前在飞行中的任务数。关键入口是 internal_try_put_taskif (my_concurrency my_max_concurrency) { my_concurrency; graph_task* new_task create_body_task(...); // 立即生成任务 ... } else if ( my_queue my_queue-push(...) ) { // 并发已满消息进入缓冲队列 op-bypass_t SUCCESSFULLY_ENQUEUED; ... }也就是说当m已经有一个任务在飞行my_concurrency 1达到上限时后续到达的消息会被缓冲到节点内部的输入队列function_input_queue定义于 _flow_graph_node_impl.h中而不是立刻生成任务。这正是用户指南指出的行为如果节点无法立即生成任务来处理消息消息将被缓冲在节点中当并发度允许时再为下一条缓冲消息生成任务。一个任务结束触发下一个任务的生成当前任务执行完毕后节点需要通过一次归还并发额度的操作来唤醒缓冲中的下一条消息。在 handle_operations 的app_body_bypass分支中可以看到case app_body_bypass: { --my_concurrency; if (my_concurrency my_max_concurrency) tmp-bypass_t perform_queued_requests(); // 从队列取消息并创建新任务 ... }perform_queued_requestsL202-L226会检查输入队列若非空就my_concurrency并调用create_body_task为队首消息生成任务、随后pop出队。于是m呈现出前一个任务完成后立即开始下一个的接力式串行执行。serial 与任意数值并发度上述机制对任意正整数并发度都成立把1换成4或8就意味着节点最多允许 4 或 8 个任务同时在飞行超出部分进入缓冲队列。oneTBB 预定义了最常用的两个枚举值flow_graph.h值语义unlimited 0消息到达即生成任务不限制在飞行任务数serial 1同一时刻至多 1 个任务执行函数体消息按序缓冲处理任意正整数 N同一时刻至多 N 个任务执行函数体超出部分缓冲可以这样记忆并发度决定的是任务生成的节奏而不是线程的数量。n生成三个任务可能并发执行m生成的任务被刻意限速为一次一个因此在 m 内部永远看不到并发。wait_for_all阻塞的是调用线程不是所有线程当所有λn与λm任务都完成后g.wait_for_all()返回示例程序结束。这一小节的三个要点值得展开一、try_put 不阻塞时间线中三次n.try_put(...)调用本身在图中占据的是操作栈一侧的薄层——它们快速返回。原因如 internal_try_put_task 所示消息要么立即被用于创建任务要么被推入缓冲队列并返回SUCCESSFULLY_ENQUEUED标记调用线程不会被卡住等待任务执行完毕。用户指南在Mapping Nodes to Tasks一节的note中明确指出流图中所有执行都是异步的——try_put在立即生成任务或缓冲消息后就把控制权交还调用线程任务执行函数体后再把结果投递给后继节点只有wait_for_all会也应该阻塞。二、wait_for_all 的调用线程会帮忙干活wait_for_all并不是让调用线程空转自旋spinning idly。oneTBB 的实现让阻塞中的调用线程转而从任务池中偷取并执行其他任务。源码注释与实现都明确这一点graph::wait_for_all 的文档注释写着The waiting thread will go off and steal work while it is blocked in the wait_for_all其内部在图的 task_arena 中执行d1::wait(...)等待图的根task_group_context。这对性能有实际意义即使你的流图只属于某一个线程wait_for_all期间该线程也不会闲着而是参与执行 oneTBB 工作池中的其他任务包括本图和其他图中的任务帮助维持线程池的整体吞吐。这是 oneTBB 工作窃取work-stealing调度与图执行结合的关键体验。三、wait_for_all 等待的是图中没有任务在跑wait_for_all的语义是阻塞直到图中不再有正在执行的任务时间线图中n的全部并发任务与m的全部串行任务均完成。配合 Graph_Object 一节的知识graph对象是所有节点的集合句柄调用它的wait_for_all等价于等待与整个图相关的全部活动结束。同一个graph对象还承担reset重置全部节点状态与cancel取消所有节点执行等整图操作。线程数不足时的时间线正文时间线展示的是线程充足的理想情形n的三个任务、以及m的接力任务都能即时找到线程执行。但需要明确该图的适用前提——如果可用线程更少部分已生成的任务只能等待直到某条线程空闲出来。例如单线程系统上n的λn(1)、λn(2)、λn(3)与m的λm(1..3)会被压缩为一条完全串行的执行链两线程系统上可能同时存在一个λn与一个λm在跑其余任务排队。无论哪种情况任务与消息的数量关系不变n侧来一条消息建一个任务m侧并发额度归还一次处理一条缓冲消息。变的只是任务在时间轴上的排布密度。在真实测试中印证unlimited 节点确实即来即建任务仓库的流图测试 test_flow_graph.cpp 中大量使用了tbb::flow::unlimited构造各类节点例如第 220 行tbb::flow::function_node int f_n(my_graph, tbb::flow::unlimited, function_bodyint(midway_arena));该测试还展示了make_edge(join_node, function_node)、make_edge(output_portN(split_node), function_node)等多节点组合验证了前驱输出自动投递给后继的边语义。你可以把unlimited换成serial或具体数值观察并发行为差异——这与本文从源码推导的记账机制完全一致。小结把本文的结论压缩成一张速查表观察点行为源码依据try_put返回立即生成任务或缓冲消息不阻塞调用线程internal_try_put_taskunlimited节点消息到达即创建任务不限制在飞行数try_put_task_impl并发度 N 节点在飞行任务数达到 N 后缓冲消息任务完成时再取队首生成新任务app_body_bypass / perform_queued_requests任务执行与转发函数体输出经successors().try_put_task(v)自动投递给后继节点apply_body_impl_bypasswait_for_all阻塞至图中无活动任务等待期间调用线程参与执行任务池中的其他任务graph::wait_for_all线程池任务数量与线程数量解耦只有线程池中的可用线程真正执行任务Nodes 一节理解节点 → 任务的映射之后你就掌握了流图性能分析的第一个基本工具看到try_put不要问它会不会阻塞要问它生成了几个任务、这些任务能并行到什么程度、消息在哪里被缓冲。沿着这个思路下一步可以继续阅读用户指南中的 Flow_Graph_Message_Passing_Protocol消息传递协议与 Flow_Graph_Buffering_in_Nodes节点缓冲把任务调度与消息语义两块拼图合拢。赞分享并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载相关推荐LDDC歌词下载神器5分钟掌握多平台精准歌词获取全攻略LDDC歌词下载神器5分钟掌握多平台精准歌词获取全攻略 还在为找不到心爱歌曲的歌词而烦恼吗LDDC歌词下载工具作为一款完全免费的精准歌词获取神器支持QQ音桌面应用从单节点到跨节点C线程池如何构建分布式任务调度系统从单节点到跨节点C线程池如何构建分布式任务调度系统 你是否遇到过这样的困境单服务器的线程池ThreadPool处理能力已达极限但业务仍需处理成百上并发编程au-automatic分布式任务调度多节点协同生成方案au automatic分布式任务调度多节点协同生成方案 引言分布式任务调度的必要性与挑战 在大规模AI生成任务场景中单节点计算资源往往成为瓶颈。au aAI 应用媒体生成计算机视觉后端上一篇抖音下载神器完整指南教你如何轻松保存无水印视频和直播内容下一篇番茄小说下载器终极指南3步实现小说离线永久保存创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表