资讯动态

oneTBB parallel_pipeline 实战解析:用 Square 示例实现文本数字平方的流水线并行

发布时间:2026/10/9 1:26:26 来源:尧图企业网站定制
并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载本文以 oneTBB 仓库中的 Square 示例 为主线完整讲解parallel_pipeline的构建、运行与三阶段过滤器读取、转换、写出设计并深入源码揭示其工作机理。读完本文你将掌握parallel_pipeline的典型用法如何用make_filter组装流水线、如何用 token 数控制并行度、以及输入过滤器如何通过flow_control::stop()优雅终止流水线。一、示例简介一个文本数字平方器Square 示例是一个典型的流水线程序读取一个包含十进制整数文本的文件把其中的每个数字替换为它的平方再写回新的输出文件。例如输入中的3会被改写为912改写为144而数字之间的空白、标点等非数字字符原样保留。示例位于仓库的examples/parallel_pipeline/square/目录包含三个源文件与一个构建脚本square.cpp主程序定义三个过滤器并驱动流水线gen_input.cpp输入文件不存在时自动生成约 100 万个测试数字CMakeLists.txtCMake 构建脚本。之所以适合用流水线模型是因为整个处理可以拆成三个阶段顺序读取文件块 → 并行地平方每块中的数字 → 按原顺序写回输出。文件 I/O 本身是顺序的但中间的数字转换只依赖块内数据、互不干扰可以充分并行——这正是parallel_pipeline的用武之地。二、构建与运行2.1 用 CMake 构建参照示例目录内的 CMakeLists.txt构建步骤如下其中path_to_example指向示例目录即examples/parallel_pipeline/squarecmake path_to_example cmake --build .CMake 脚本通过include(../../common/cmake/common.cmake)复用公共配置见 examples/common/cmake/common.cmake并通过find_package(TBB ...)链接TBB::tbb与Threads::Threads两个目标同时把square可执行文件由gen_input.cpp与square.cpp共同编译而成。2.2 预定义 make 目标构建完成后示例提供以下便捷执行目标由 examples/common/cmake/common.cmake 中的add_execution_target宏生成参数定义在 CMakeLists.txt 末尾目标作用实际执行参数make run_square以预定义参数运行示例0 input.txt output.txtmake perf_run_square以推荐参数运行用于测量 oneTBB 性能auto input.txt output.txt silentmake light_test_square以精简参数快速跑通缩短执行时间由 examples/CMakeLists.txt 的聚合框架决定若某示例未显式定义light_test_name目标则回退到run_nameperf_run_square中auto表示线程数取平台默认值silent让程序只输出耗时便于基准测量。三、命令行参数完整用法程序的完整用法如下[...]表示可选square [n-of-threadsvalue] [input-filevalue] [output-filevalue] [max-slice-sizevalue] [silent] [-h] [n-of-threads [input-file [output-file [max-slice-size]]]]参数说明参数含义默认值 / 取值-h打印命令行选项帮助隐式提供见 square.cpp 中utility::parse_cli_arguments的调用n-of-threads使用的线程数可以是形如low[:high]的范围low与可选high为非负整数或auto平台默认线程数默认 0表示先串行再全并行运行input-file输入文件名默认input.txtoutput-file输出文件名默认output.txtmax-slice-size单个切片slice的最大字符数默认 4000对应源码中的MAX_CHAR_PER_INPUT_SLICE见 square.cppsilent除耗时外不输出任何内容布尔开关参数解析由公共工具 examples/common/utility/utility.hpp 中的utility::parse_cli_arguments完成支持keyvalue与位置参数两种写法。n-of-threads的auto值取自 examples/common/utility/get_default_num_threads.hpp其内部调用oneapi::tbb::this_task_arena::max_concurrency()。主程序还会调用generate_if_needed()定义于 gen_input.cpp当输入文件不存在时自动生成 100 万个 0–9999 范围内的整数写入文件保证示例开箱即用。四、流水线结构三个过滤器串成一条装配线parallel_pipeline把数据处理建模成传统制造业的装配线数据以token形式流经一串过滤器filter每个过滤器对数据做一部分加工。示例的流水线由三个过滤器串联而成完整调用见 square.cpponeapi::tbb::parallel_pipeline( nthreads * 4, oneapi::tbb::make_filtervoid, TextSlice*(oneapi::tbb::filter_mode::serial_in_order, MyInputFunc(input_file)) oneapi::tbb::make_filterTextSlice*, TextSlice*(oneapi::tbb::filter_mode::parallel, MyTransformFunc()) oneapi::tbb::make_filterTextSlice*, void(oneapi::tbb::filter_mode::serial_in_order, MyOutputFunc(output_file)));三个过滤器各司其职阶段过滤器类型模式职责输入make_filtervoid, TextSlice*serial_in_order顺序从文件读入文本块包装为TextSlice转换make_filterTextSlice*, TextSlice*parallel并行地把块内每个十进制数替换为其平方输出make_filterTextSlice*, voidserial_in_order按原顺序把处理后的块写回输出文件关键点make_filterInputType, OutputType(mode, functor)创建过滤器InputType是流入值类型OutputType是流出值类型。第一个过滤器的输入为void最后一个过滤器的输出为void详见 include/oneapi/tbb/parallel_pipeline.h。operator串联过滤器时前一过滤器的OutputType必须与后一过滤器的InputType一致。在本例中均为TextSlice*且通过指针传递切片避免复制整个文本块的开销见 include/oneapi/tbb/parallel_pipeline.h。filter_mode决定并发语义共三种见 include/oneapi/tbb/parallel_pipeline.hparallel可同时处理多个 token且不保证顺序serial_in_order一次只处理一个 token且所有serial_in_order过滤器保持相同的处理顺序serial_out_of_order一次只处理一个 token但不保证顺序。本示例中输入过滤器读取顺序文件、输出过滤器必须保持文件原始顺序因此两者必须是serial_in_order中间转换只依赖块内数据故声明为parallel。若某 token 在输出过滤器处乱序到达流水线会自动延迟调用输出函数直到其前驱处理完毕——这正是serial_in_order的保序保证。五、token 数量并行度的节流阀parallel_pipeline的第一个参数nthreads * 4是允许同时在流水线中流通的最大 token 数即在途 token 上限。示例注释明确指出每个线程需要 2–4 个在途 token 才能让所有线程保持忙碌。为什么必须限制 token 数设想 token 无上限中间的不保序过滤器会不断吞入新 token而输出过滤器跟不上消费速度导致中间过滤器无节制地占用内存等资源。一旦达到上限流水线便不再在输入过滤器处创建新 token直到某个 token 在输出端被销毁。这一机制在底层由 src/tbb/parallel_pipeline.cpp 的pipeline构造器实现其中input_tokens(Token(max_token))初始化了 token 计数并断言max_token 0。此外还有一层并发约束parallel_pipeline本身在调用线程所在的 task arena 中执行实际并发度还受 arena 中可用线程数的限制本示例通过oneapi::tbb::global_control的max_allowed_parallelism控制。六、数据载体 TextSlice头尾分离的变长块为了摊薄调度开销示例按约 4000 字符的块处理文本。每个块用TextSlice表示见 square.cpp其设计要点是对象只保存头部两个指针字符数据紧随对象之后分配因此必须通过allocate/free分配和释放不能按普通对象析构logical_end指向已填充内容的末尾physical_end指向缓冲区末尾二者之间是可用空间avail()allocate(max_size)用oneapi::tbb::tbb_allocatorchar()分配sizeof(TextSlice) max_size 1字节多出的 1 字节留给终止\0供strtol安全解析。七、三个过滤器的内部实现7.1 输出过滤器最简单的一环MyOutputFunc见 square.cpp只做两件事用fwrite把TextSlice的字符写入文件然后调用out-free()释放切片。若写入字节数与切片长度不符则向 stderr 报告错误并退出。7.2 转换过滤器并行求平方MyTransformFunc见 square.cpp逐字符扫描输入切片先给切片末尾补\0保证strtol在数字恰好位于块尾时也能正确解析原样拷贝非数字字符遇到数字时用strtol解析为long计算平方后用sprintf写回。注释说明了为何无需溢出检查输出缓冲区按输入的两倍分配2 * MAX_CHAR_PER_INPUT_SLICE而非负整数n的平方n²的十进制位数不可能超过n的两倍。处理完毕后释放输入切片返回新切片。由于它只操作块内数据多个调用可并发执行因此声明为parallel模式。7.3 输入过滤器最复杂负责边界与终止MyInputFunc见 square.cpp是最复杂的一环它必须解决两个问题1数字不能跨越切片边界。若某次读取恰好填满缓冲区块尾可能截断一个数字例如读到12的一半。此时算法从块尾向前回退到最后一个非数字字符把半截数字搬到下一个切片开头保证数字的完整性。若回退到块首仍未找到非数字则说明该数字太大、放不进单个缓冲区此时触发断言。2何时宣告输入结束。当fread返回 0 且当前切片为空时说明文件已读完此时调用fc.stop()并返回nullptr。flow_control参数是第一个过滤器输入过滤器独有的——任何作为流水线首端的 functor 都必须接受它并用stop()通知流水线停止生成新 token对应 include/oneapi/tbb/parallel_pipeline.h 导出的flow_control底层flow_control::stop()语义见 src/tbb/parallel_pipeline.cpp 附近注释。另外由于 functor 在构造过滤器链与启动流水线时都会被复制MyInputFunc显式定义了拷贝构造函数见 square.cpp确保多个副本共享同一个输入文件句柄。八、两个必须遵守的约束functor 的operator()必须声明为const。原因在于过滤器体对象可能被复制到其他线程若operator()修改了成员修改发生在原对象还是副本上是不确定的。parallel_pipeline通过要求const operator()从编译期杜绝这类数据竞争官方指南 doc/main/tbb_userguide/Working_on_the_Assembly_Line_pipeline.rst 有专门 CAUTION 说明。首端过滤器的operator()必须接受flow_control并在输入耗尽时调用stop()否则流水线无法正常结束。九、从源码看并行度与实现parallel_pipeline的公开接口定义在 include/oneapi/tbb/parallel_pipeline.h三个模板参数版本的make_filter支持 C17 推导指引filter(filter_mode, Body)允许省略模板实参支持两种调用形式把过滤链预先组装成filtervoid,void再传入或把多个过滤器直接作为可变参数传入内部自动用operator连接还提供带task_group_context的重载用于把流水线纳入用户自定义的任务组上下文便于统一取消与异常传播。底层运行时实现在 src/tbb/parallel_pipeline.cppparallel_pipeline构造pipeline对象用Token(max_token)管理在途 token 计数并在 token 数降为 0即所有输入消费完毕时结束执行。十、扩展阅读官方用户指南对流水线模式有完整讲解Working on the Assembly Line: parallel_pipeline内容与本示例的代码一一对应若流水线中存在环形依赖可阅读 Using Circular Buffers吞吐量调优与非线性流水线分别见 Throughput of pipeline 与 Non-Linear Pipelines更多可运行示例位于仓库 examples 目录其中graph/下另有以 flow graph 实现的并行处理示例可作为对照学习。赞分享并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载相关推荐django CMS 多站点Multi-Site部署实战指南SITE_ID 经典方案与 5.1 请求级站点解析django CMS 多站点Multi Site部署实战指南SITE_ID 经典方案与 5.1 请求级站点解析 django CMS 基于 Django并发编程高性能计算oneTBB parallel_pipeline 流水线并行模式实战从装配线模型到源码级原理oneTBB parallel_pipeline 流水线并行模式实战从装配线模型到源码级原理 本指南以 oneAPI Threading Building B并发编程高性能计算oneTBB 循环与流水线算法全景指南parallel_for、parallel_reduce、parallel_for_each 与 parallel_pipeline 实战详解oneTBB 循环与流水线算法全景指南parallel_for、parallel_reduce、parallel_for_each 与 parallel_pi并发编程高性能计算上一篇QtScrcpy免费开源的Android投屏控制终极指南3个核心功能让你轻松上手下一篇pwndbg 的 hi 命令堆地址归属快速定位与 malloc_chunk 逆向解析实战创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑