ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Arrow C++ 异步执行机制图解:CPU 与 I/O 线程池如何通过 Future 与 Continuation 协同调度

Apache Arrow C++ 异步执行机制图解:CPU 与 I/O 线程池如何通过 Future 与 Continuation 协同调度 Apache Arrow C 异步执行机制图解CPU 与 I/O 线程池如何通过 Future 与 Continuation 协同调度【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrowApache Arrow C 的 Acero 执行引擎通过CPU 线程池 I/O 线程池 Future/Continuation的组合方式实现异步执行从而在 I/O 等待期间不浪费 CPU 核心。本文以仓库中docs/source/developers/cpp/img/async.md这张官方序列图为骨架逐步拆解异步任务从提交、阻塞、让出线程到续跑完成的完整生命周期并对照future.h、thread_pool.h、async_util.h与scan_node.cc等源码帮助读者理解 Arrow 底层异步机制的设计动机与实现细节为阅读 Acero 执行计划源码或自行开发自定义 ExecNode 打下基础。这张图从哪来async.md 是 async.svg 的源码在 Arrow 仓库中docs/source/developers/cpp/img/async.md并非一篇独立的教程而是用于生成async.svg架构图的 Mermaid 序列图源文件。文件开头两行注释明确写道This is the source for the async.svg diagram used in developer_guide.rst这张 SVG 图被引用在 Acero 开发者指南 acero.rst 的 Asynchronicity异步性小节配图说明为Arrow achieves asynchronous execution by combining CPU I/O thread pools也就是说这张图回答了一个核心问题Arrow 如何利用两种分工不同的线程池让长时间阻塞的 I/O 操作不再占用宝贵的 CPU 执行线程。本文接下来逐条解读图中每一步的含义并结合源码说明其背后的实现机制。序列图逐帧拆解一次异步 I/O 的完整生命周期原图是一段 22 行的 MermaidsequenceDiagram参与者有三个Thread Pool线程池、CPUCPU 执行线程和IOI/O 等待线程。下面按时间顺序拆解完整流程。第一步提交任务发起异步读取Thread Pool-CPU: Start Task CPU-IO: Read计划启动时线程池将一个任务调度到 CPU 线程上执行。该任务在执行过程中需要读取数据例如扫描文件于是向 I/O 线程发起一次异步Read请求。注意这里的读取是非阻塞的CPU 线程发出请求后并不会原地等待而是立即返回。第二步I/O 线程返回 FutureCPU 线程挂接 ContinuationIO-CPU: FutureBuffer CPU-CPU: Add Continuation CPU-Thread Pool: Finish TaskI/O 线程接受读取请求后立即向 CPU 线程返回一个FutureBuffer——一个尚未完成、代表未来某个时刻会拿到 Buffer的句柄。CPU 线程拿到 Future 后不是阻塞等待而是往这个 Future 上挂一个 Continuation续接回调然后立刻结束当前任务把 CPU 线程归还给线程池。这正是异步编程的核心思想把等待结果替换为注册回调让出线程让其他任务运行。对应到源码FutureT提供AddCallback与Then方法用于注册续接逻辑例如 future.h 中的AddCallback与 future.h 中的ThenThen会基于回调的返回值推导出一个新的Future类型实现链式续接。第三步I/O 线程阻塞等待CPU 线程池继续服务其他任务Note right of IO: Blocked on IO Thread Pool-CPU: Other Task CPU-Thread Pool: Thread Pool-CPU: Other Task CPU-Thread Pool: Thread Pool-CPU: Other Task CPU-Thread Pool:图中用Note right of IO: Blocked on IO明确标注现在只有 I/O 线程处于阻塞状态。与此同时线程池把三个Other Task依次派发给 CPU 线程CPU 线程逐个执行完毕并归还。这直观地展示了异步模型的收益一次文件读取可能耗时数毫秒如果没有异步机制等待期间这颗 CPU 核心就完全闲置而在异步模型下同样的时间窗口内可以执行多个其他计算任务。这就是 acero.rst 中所说的两种解决思路的取舍同步方案创建比核心数更多的线程容忍部分线程阻塞。实现简单但会产生线程竞争且需要精细调优。异步方案慢操作启动后调用方让出线程操作完成时再创建新任务继续处理结果。线程竞争最小但实现更复杂。由于 C 标准库缺少统一的异步 APIAcero 选择两者结合CPU 线程池每个核心一个线程任务绝不应该阻塞除轻微同步延迟外应尽可能持续占用 CPUI/O 线程池的线程则大部分时间处于空闲等待状态几乎不做 CPU 密集型工作其职责就是等数据就绪然后在 CPU 线程池上调度后续任务。第四步读取完成续接执行deactivate IO IO-IO: Read Finished IO-IO: Run Continuation IO-Thread Pool: Schedule TaskI/O 操作完成后I/O 线程标记读取完成运行之前注册好的 Continuation并把后续任务重新调度回线程池。图中deactivate IO表示 I/O 线程从阻塞状态解除。第五步结果回到 CPU 线程继续处理Thread Pool-CPU: Start Task CPU-CPU: Process Read Result线程池再次把一个新任务派发给 CPU 线程此时读取结果已经就绪CPU 线程直接执行Process Read Result——例如把读到的数据块解码、投影或传递给下游节点。至此一次完整的异步读取闭环结束。整个流程的核心可以概括为一句话发起 I/O 的 CPU 线程不会陪 I/O 一起等待而是通过 Future 挂回调、立刻让出线程I/O 完成后再以新任务的形式把控制权交还 CPU 线程池。线程池的源码实现CPU 与 I/O 两套 Executor图中Thread Pool / CPU / IO三个参与者对应源码中 Arrow 的全局线程池体系。在 thread_pool.h 中Arrow 提供了查询与调整全局 CPU 线程池容量的接口GetCpuThreadPoolCapacity()返回 CPU 线程池容量一个理想值不一定是某一时刻的实际线程数SetCpuThreadPoolCapacity(int threads)设置 CPU 线程池的工作线程数。线程池本身基于Executor抽象thread_pool.h其关键能力包括Spawn提交一个即发即忘fire-and-forget任务任务无返回值适合派发一次性工作Submit提交可调用对象并返回一个 Future用于获取执行结果Transfer(Future)把 Future 迁移到指定执行器上——这是异步模型中保证续接回调跑在正确的线程池的关键 API。Transfer的注释非常直白地说明了设计意图thread_pool.h当 I/O 任务完成一个 Future 时该 Future 的 Continuation 默认会在调用MarkFinished的那个线程即 I/O 线程上执行为了让 CPU 密集型工作不占用 I/O 线程池I/O 任务应当把 FutureTransfer到 CPU 执行器后再返回。这与序列图中IO 完成读取 → Schedule Task 回到线程池的步骤完全对应。从源码结构看Arrow 由此形成了明确的分工约定I/O 线程只负责等待与唤醒CPU 线程负责一切实际计算。任务调度的进阶工具AsyncTaskScheduler 与节流调度器序列图展示的是单次异步读取的微观流程而 Acero 执行引擎在宏观层面还需要对大量并发任务进行调度与控制。Acero 开发者指南的 Thread Pools and Schedulers 一节acero.rst指出CPU 与 I/O 线程池本身只提供 FIFO 任务队列Acero 在此基础上使用AsyncTaskSchedulerasync_util.h来获得三方面额外能力节流ThrottlingThrottledAsyncTaskSchedulerasync_util.h为每个任务关联一个代价cost只有当当前并发代价总和不超过max_concurrent_cost时才提交任务否则进入队列。该调度器还可手动Pause()/Resume()暂停后即使有空间也不会提交队列中的任务这与 ExecNode 的背压backpressure机制直接挂钩。Acero 中write节点就使用大小为 1 的节流避免重入地调用 dataset writer因为 writer 内部有自己的调度逻辑。优先级Priority可以为节流队列配置自定义Queue控制排队任务的提交顺序。任务在节流未满时立即提交无论优先级只有节流满时优先级才起作用。scan节点用它限制并发读请求数量并尽量按数据集顺序读取。任务组Task Group跟踪一组任务的完成情况全部完成后执行一个收尾任务适用于 fork-join 型问题。write节点用任务组在某个文件的全部写任务完成后关闭该文件。ThrottledAsyncTaskScheduler::Make的签名与语义在 async_util.h 中有完整说明max_concurrent_cost是任意时刻允许运行的最大代价单个任务代价超过上限时会被压到上限值以保证任务仍可运行默认使用 FIFO 队列需要优先级时可传入自定义Queue。源码实例scan 节点如何运用节流调度器序列图中的Read步骤在真实 Acero 计划里最常见于扫描节点。从 scan_node.cc 的StartProducing可以看到实际用法Status StartProducing() override { NoteStartProducing(ToStringExtra()); batches_throttle_ util::ThrottledAsyncTaskScheduler::Make( plan_-query_context()-async_scheduler(), options_.target_bytes_readahead 1); plan_-query_context()-async_scheduler()-AddSimpleTask( [this] { return GetFragments(options_.dataset.get(), options_.filter) .Then(this { ScanFragments(frag_gen); }); }, ScanNode::ListDataset::GetFragmentssv); return Status::OK(); }这段代码印证了序列图的多个环节使用GetFragments(...).Then(...)把获取数据集分片列表的续接逻辑挂到 Future 上——对应图中的Add Continuation通过ThrottledAsyncTaskScheduler::Make创建基于target_bytes_readahead目标预读字节数的节流调度器控制并发 I/O 量——对应图中 I/O 线程上受限的并发读取在ScanFragments中还用MakeThrottledAsyncTaskGroup结合fragment_readahead 1限制并发分片任务数scan_node.cc分片处理完后再通过回调output_-InputFinished(...)通知输出端。另外该节点的PauseProducing/ResumeProducing目前仍是 TODO 占位scan_node.cc但acero.rst中描述了预期的行为暂停时冻结调度器节流队列中尚未提交的任务不再提交而已在进行中的 I/O 会继续背压无法立即生效恢复时解冻调度器。可见节流调度器正是为支撑 ExecNode 背压语义而设计的。异步机制与 ExecNode 生命周期的关系理解序列图后再回看 Acero 的 ExecNode 生命周期会清晰很多。Acero 的执行模型是 push-based每个节点通过InputReceived把数据推给下游。而 acero.rst 特别强调大多数 Acero 节点不需要关心异步——它们完全是同步的不产生任务。只有两类节点深度依赖异步机制Source 节点如scan、table_source在StartProducing中调度读取任务是计划中任务的主要生产者也是这张异步序列图的主角Pipeline Breaker 节点如order_by、write需要积累全部输入后才工作通常在积累完成后才调度后续任务并用AsyncTaskScheduler跟踪任务完成情况。PauseProducing/ResumeProducing则负责背压当下游如 SinkNode 的队列或 write 节点的写入队列积压过满时上游调用PauseProducing暂停生产队列排空后再ResumeProducingacero.rst。write节点在InputReceived中把 batch 加入写入队列若队列已满则获得一个未完成的 Future并给该 Future 挂一个恢复 Continuation——这与序列图中挂回调 → 让出 → 完成后再调度的模式如出一辙。设计哲学小结序列图背后体现的是 Acero 的三条核心设计哲学acero.rstMake Tasks not Threads需要并行时使用线程池任务而非专用线程以降低线程竞争与上下文切换、避免死锁失败时任务自动取消、简化性能剖析并支持无线程环境如 emscriptenDont Block on CPU Threads耗时且不占用 CPU 的活动磁盘读取、网络 I/O、等待外部库必须使用异步工具绝不能让 CPU 线程阻塞——这正是 async.svg 序列图要传达的核心信息Task per Pipeline任务尽量贯穿整条流水线减少中间节点间的数据传递与缓存失效而异步 I/O 的引入保证了这种流水线模型在等待数据时依然能保持 CPU 利用率。如果希望在阅读源码时对照这张图可以直接查看 async.md 中的 Mermaid 源码也可以阅读 acero.rst 中对应的 Asynchronicity 章节再结合 future.h、thread_pool.h、async_util.h 与 scan_node.cc 的注释与实现逐行印证即可完整掌握 Arrow C 的异步执行全貌。【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表