news 2026/9/10 11:38:41

Folly Executors 与线程池深度指南:CPUThreadPoolExecutor 与 IOThreadPoolExecutor 的原理、选型与实践

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Folly Executors 与线程池深度指南:CPUThreadPoolExecutor 与 IOThreadPoolExecutor 的原理、选型与实践

Folly Executors 与线程池深度指南:CPUThreadPoolExecutor 与 IOThreadPoolExecutor 的原理、选型与实践

【免费下载链接】follyAn open-source C++ library developed and used at Facebook.项目地址: https://gitcode.com/GitHub_Trending/fol/folly

Folly 在 folly/docs/Executors.md 中系统讲解了其线程池(Thread Pool)与 Executor 体系的设计动机、实现原理和使用方式。本文以该文档为骨架,结合 folly/executors 目录下的源码实现,深入剖析 Folly 为何要自研两套线程池(CPU 与 IO 分离)、它们各自的数据结构与调度细节,以及如何借助全局 Executor、Observers 与 PoolStats 在生产环境中高效运行并发代码。读完本文,你将掌握CPUThreadPoolExecutorIOThreadPoolExecutorThreadPoolExecutor基类的核心机制,并能在自己的代码中正确选型与配置。

快速上手:从全局 Executor 开始

Folly 提供两个具体的线程池实现:IOThreadPoolExecutorCPUThreadPoolExecutor,并将它们作为完整异步框架(folly/futures)的一部分内置。大多数场景下,你并不需要手动构造线程池,而是直接获取进程级的全局 Executor,再与 Future 组合使用。

最简单、最常见的用法是将 Future 延续(continuation)调度到 CPU 线程池上执行:

auto f = someFutureFunction().via(getCPUExecutor()).then(...);

如果需要在 Thrift / memcache 客户端上发起调用,则需要先拿到一个事件循环(EventBase),此时从 IO 线程池获取getEventBase(),再通过via(getCPUExecutor())回到 CPU 线程池处理结果:

auto f = getClient(getIOExecutor()->getEventBase())->callSomeFunction(args...) .via(getCPUExecutor()) .then([](Result r) { /* do something with result */ });
  • getCPUExecutor()返回全局 CPU 线程池(folly/executors/GlobalExecutor.h);
  • getIOExecutor()返回全局 IO 线程池,可用->getEventBase()取出一个按 round-robin 方式选中的EventBase,直接在上面调度 IO 工作。

全局 Executor 的线程数与新旧 API

全局 Executor 的线程数量可以通过 gflags 显式配置:

  • folly_global_cpu_executor_threads:全局 CPU 线程池创建的线程数;
  • folly_global_io_executor_threads:全局 IO 线程池创建的线程数。

两者默认值均为 0。在 folly/executors/GlobalExecutor.cpp 中可以看到,当 flag 为 0 时,线程数会回退为folly::available_concurrency(),即机器可用的核数:

size_t nthreads = FLAGS_folly_global_cpu_executor_threads; nthreads = nthreads ? nthreads : folly::available_concurrency(); return new std::shared_ptr<ImmutableGlobalCPUExecutor>( new ImmutableGlobalCPUExecutor( nthreads, std::make_shared<NamedThreadFactory>("GlobalCPUThreadPool")));

需要说明的是,文档中示例使用的getCPUExecutor()/getIOExecutor()在当前版本中已被标记为 deprecated。从 folly/executors/GlobalExecutor.h 的注释可以看出,官方推荐:

  • 使用getGlobalCPUExecutor()/getGlobalIOExecutor()获取不可变的全局 Executor(返回KeepAlive,可安全配合 Future / coroutine 使用,保证前向进度);
  • 若确有替换全局 Executor 的需求,使用getUnsafeMutableGlobalCPUExecutor()/setUnsafeMutableGlobalCPUExecutor()等 "UnsafeMutable" 系列接口;
  • getGlobalCPUExecutorWeakRef()返回弱 KeepAlive:不阻止全局 Executor 在关闭时析构,适合"尽力而为"的后台任务,但不应用于 Future / coroutine 延续,因为延续依赖前向进度保证,弱引用可能导致死锁。

为什么不用 C++11 的 std::launch?

C++11 的std::launch只有两种模式:asyncdeferred,在生产系统中两者都不是理想选择:

  • async:每次 launch 都会无限制地新起一个线程,线程数量完全不受控;
  • deferred:任务被延迟到"真正需要结果时"才执行,且届时在当前线程里同步地运行,阻塞调用方。

Folly 的线程池则不同:任务总是尽可能早地被调度执行,同时对最大任务数 / 线程数设有上限,因此永远不会使用超出需要的线程数。这种"按需生长、有界并发"的模型正是生产级服务所需要的。

为什么需要自研线程池?

文档给出的理由是:当时(C++11 时代)现成的线程池实现都不完整——基于 pipe 的线程池太慢;而多个较早的实现不支持std::function,无法与 Future/回调风格的代码无缝协作。因此 Folly 需要自己提供一套面向高并发、低延迟、且与std::function和异步框架深度集成的线程池实现。

为什么需要两种不同类型的线程池?

这是理解 Folly 线程池架构最关键的问题,核心在于IO 事件循环与公平队列在操作系统原语层面互相排斥

  1. epoll 需要 fdevent_fd是最新的通知机制,但它有个副作用——一个活跃的 fd 会触发所有正在等待它的 epoll 循环(惊群效应,thundering herd)。因此,如果你想要一个公平队列(一个总队列 vs. 每工作线程一个队列),就需要借助信号量(semaphore)来实现公平唤醒;
  2. 信号量进不了 epoll 循环:semaphore 无法被放入 epoll 等待集合中,所以基于信号量的公平队列与 IO 事件循环不兼容;
  3. IO 与 CPU 本就该分离:即便技术上可行,通常也应当把 IO 密集与 CPU 密集的工作分开,以便在 IO 路径上获得更强的尾延迟(tail latency)保证。

正是基于上述原因,Folly 提供了两套定位不同的线程池:IOThreadPoolExecutor(面向事件循环与 IO)和CPUThreadPoolExecutor(面向计算密集任务)。

IOThreadPoolExecutor:event_fd + 每线程 NotificationQueue

IOThreadPoolExecutor是一个面向 IO 密集型任务的线程池,folly/executors/IOThreadPoolExecutor.h 的类注释完整概括了它的设计要点:

  • 使用event_fd进行通知并唤醒 epoll 循环
  • 每个线程(每个 epoll)对应一个队列,具体是NotificationQueueIOThread结构中持有自己的EventBase*,任务通过ioThread->eventBase->runInEventBaseThread(...)投递(见 folly/executors/IOThreadPoolExecutor.cpp);
  • 无谓系统调用被消除:如果目标线程已经在运行、并没有阻塞在 epoll 上等待,那么只需要把新任务放进它的队列即可,无需额外的 syscall 去唤醒事件循环;
  • 空闲线程回收内存:如果某个线程等待超过数秒,它的栈会被madvise掉。实现上由MemoryIdlerTimeout完成——它是AsyncTimeout+EventBase::LoopCallback的合体,事件循环空闲一段时间后会调用MemoryIdler::flushLocalMallocCaches()MemoryIdler::unmapUnusedStack(...)(见 folly/executors/IOThreadPoolExecutor.cpp)。不过文档也指出:当前任务在队列间是 round-robin 调度的,所以除非系统完全没有工作,否则该优化效果有限;
  • getEventBase()按 round-robin 选择:调用方可以直接拿到一个EventBase在上面调度 IO 工作;
  • 队列几乎无竞争:由于每线程一个队列,队列上的竞争极小,因此用"自旋锁 +std::deque"这种轻量结构即可承载任务,且没有最大队列长度限制
  • 默认每核一线程:默认情况下 IO 线程数与 CPU 核数一致——只要这些线程不阻塞,配置比核数更多的 IO 线程通常没有意义。

构造与选项

explicit IOThreadPoolExecutor( size_t numThreads, std::shared_ptr<ThreadFactory> threadFactory = std::make_shared<NamedThreadFactory>("IOThreadPool"), folly::EventBaseManager* ebm = folly::EventBaseManager::get(), Options options = Options());

Options支持三个配置项(folly/executors/IOThreadPoolExecutor.h):

配置项默认值说明
waitForAllfalse析构/stop()时是否等待事件循环完全退出
enableThreadIdCollectionfalse是否显式开启线程 ID 收集(返回WorkerProvider
maxReadAtOnce取 flagfolly_iothreadpoolexecutor_max_read_at_once(默认-1,即不设限)事件循环中每次循环最多读取的事件数

其中folly_iothreadpoolexecutor_max_read_at_oncedynamic_iothreadpoolexecutor(默认true,即 IO 线程池动态创建线程,最小线程数从 0 开始)两个 gflags 定义在 folly/executors/IOThreadPoolExecutor.cpp。

stop() 的行为差异

文档与源码都特别强调了一个容易踩坑的点:对于 IOThreadPoolExecutor,stop()表现得像join()。因为未完成的任务属于事件循环(EventBase),它们会在 EventBase 析构时被继续执行,所以stop()会等待这些任务完成。这与 CPUThreadPoolExecutor 的"尽力而为"停止语义不同,详见 folly/executors/IOThreadPoolExecutor.h。

CPUThreadPoolExecutor:LifoSem + MPMC 单队列

CPUThreadPoolExecutor是面向 CPU 密集任务的线程池,folly/executors/CPUThreadPoolExecutor.h 的类注释完整说明了其设计:

  • 单一队列,后端是folly::LifoSem+folly::MPMCQueue。由于全局只有一个队列,所有工作线程与所有生产者线程都会命中同一条队列,竞争可能相当高——而 MPMC 队列(多生产者多消费者无锁队列)恰好在这种场景下表现出色;
  • MPMC 队列决定了存在最大队列长度:有界队列在满时的行为由QueueBehaviorIfFull决定。以 folly/executors/task_queue/LifoSemMPMCQueue.h 为例,THROW模式会在队列满时抛出QueueFullExceptionBLOCK模式则阻塞写端;
  • LifoSem 按 LIFO 顺序唤醒线程:即始终只保持"恰好够用"的少数线程在运行,并尽量复用同一批线程以获得更好的 cache locality;其余线程被挂起,直到出现工作尖峰时才被唤醒;
  • 空闲线程栈回收:所有 Folly 的BlockingQueue实现都基于LifoSemThrottledLifoSem,长期不活跃的线程,其栈会被madvise掉;
  • stop()会在退出时完成所有未完成任务
  • 支持优先级:优先级通过多条队列实现——每个工作线程总是先检查最高优先级的队列。但线程本身并不设置 OS 优先级(pthreads 线程优先级在实践中的表现不佳),因此一连串长时间运行的低优先级任务仍可能占满所有线程。

队列工厂:默认、LIFO 与节流 LIFO

CPUThreadPoolExecutor提供了一组静态工厂方法用于构造任务队列(folly/executors/CPUThreadPoolExecutor.cpp):

工厂方法后端特点
makeDefaultQueue()由 flagfolly_cputhreadpoolexecutor_use_throttled_lifo_sem决定默认走ThrottledLifoSem(若 flag 开启),否则走LifoSem
makeLifoSemQueue()UnboundedBlockingQueue<CPUTask, LifoSem>无界、LIFO 唤醒
makeThrottledLifoSemQueue(wakeUpInterval)UnboundedBlockingQueue<CPUTask, ThrottledLifoSem>无界、节流唤醒,可配置唤醒间隔
makeDefaultPriorityQueue(numPriorities)/makeLifoSemPriorityQueue/makeThrottledLifoSemPriorityQueuePriorityUnboundedBlockingQueue<...>对应上述语义的多优先级队列版本

使用有界优先队列时,可用PriorityLifoSemMPMCQueue<CPUTask>(numPriorities, maxQueueSize)(在maxQueueSize构造重载中使用)。

死锁警告

文档对应的源码注释(folly/executors/CPUThreadPoolExecutor.h)给出了一个重要的实践警告:如果使用有界队列QueueBehaviorIfFull::BLOCK),且线程池中的任务还会继续向该线程池提交任务,一旦队列变满就可能死锁;多个使用阻塞队列的线程池之间存在环形依赖时同样可能死锁。规避方式有二:

  • 只用无界队列(默认且推荐的做法);
  • 或者只从不属于该线程池的线程中提交任务。

构造与动态线程数

CPUThreadPoolExecutor提供多种构造函数重载,包括:

  • CPUThreadPoolExecutor(size_t numThreads, Options opt = {}):最简单形态;
  • CPUThreadPoolExecutor(size_t numThreads, int8_t numPriorities, ...):多优先级版本;
  • CPUThreadPoolExecutor(size_t numThreads, int8_t numPriorities, size_t maxQueueSize, ...):多优先级 + 有界队列版本;
  • CPUThreadPoolExecutor(std::pair<size_t, size_t> numThreads, ...):显式指定(maxThreads, minThreads)

关于动态线程数,folly/executors/CPUThreadPoolExecutor.cpp 显示:当 flagdynamic_cputhreadpoolexecutor为真时,numThreads会被展开为(numThreads, 0),即最小线程数为 0、由线程池按负载动态创建/回收线程;为假时则固定线程数。构造时setNumThreads(numThreads.first)会把maxThreads_设为指定值,minThreads_会在动态模式下被置为 1(保证至少有一个可运行的线程)。

线程池可以在minThreads_maxThreads_之间动态变化实际运行线程数(由activeThreads_追踪)。任务入队后,如果无法保证已有活跃线程会处理它,就会调用ensureActiveThreads()启动新线程,直至达到maxThreads_;空闲线程的回收则交由各子类的空闲超时机制完成。相关机制见 folly/executors/ThreadPoolExecutor.h 的类注释与 folly/executors/ThreadPoolExecutor.cpp 的setNumThreads实现。

任务过期回调

基类ThreadPoolExecutor::add(func, expiration, expireCallback)提供了任务过期语义:如果func在入队后expiration时间内尚未开始执行,则执行expireCallback(见 folly/executors/ThreadPoolExecutor.h)。CPUThreadPoolExecutor还额外提供带优先级的add(func, priority, expiration, expireCallback)重载。

ThreadPoolExecutor:共享的基类逻辑

ThreadPoolExecutor是所有具体线程池的基类,包含线程的启动 / 停止 / 统计逻辑——这些逻辑与"任务具体如何被运行"是解耦的,因此被提取为公共基类(folly/executors/ThreadPoolExecutor.h)。

它提供的核心能力包括:

  • 线程生命周期管理addThreads/removeThreads/joinStoppedThreads/stopAndJoinAllThreads
  • 动态调线程数setNumThreadsnumThreads()numActiveThreads()setThreadDeathTimeout()
  • 停止与等待stop()(尽力而为,未执行任务不保证执行完)与join()
  • 批量遍历withAll(FunctionRef<void(ThreadPoolExecutor&)>)用于对所有已注册线程池执行操作,主要服务于统计导出;
  • 线程工厂setThreadFactory/getThreadFactory,默认实现是NamedThreadFactory(如"CPUThreadPool""IOThreadPool""GlobalCPUThreadPool"等命名前缀,见 folly/executors/thread_factory);
  • CPU 时间统计getUsedCpuTime()返回线程池所有线程(含已退出线程)累计的 CPU 时间,需要系统支持 per-thread CPU clock,否则返回 0,且该操作可能较昂贵。

Observers:监听线程的创建与销毁

ThreadPoolExecutor::Observer是一个观察者接口,用于监听线程的 start/stop 事件(folly/executors/ThreadPoolExecutor.h):

class Observer { public: virtual ~Observer() = default; virtual void threadStarted(ThreadHandle*) noexcept {} virtual void threadStopped(ThreadHandle*) noexcept {} virtual void threadPreviouslyStarted(ThreadHandle* h) noexcept { threadStarted(h); } virtual void threadNotYetStopped(ThreadHandle* h) noexcept { threadStopped(h); } };

它的典型用途是创建"每线程一份"的对象(如 thread-local 资源、连接池、性能计数器),同时保证在线程被动态加入或移出线程池时这些对象也能被正确地创建与清理。通过addObserver/removeObserver注册。

IOThreadPoolExecutor还扩展了IOObserver接口,额外提供registerEventBase(EventBase&)/unregisterEventBase(EventBase&)钩子(见 folly/executors/IOThreadPoolExecutor.h),便于在 IO 线程的 EventBase 创建/销毁时执行对应初始化与清理。

Stats:线程池运行指标

ThreadPoolExecutor::PoolStats结构体提供了线程池级别的统计信息(folly/executors/ThreadPoolExecutor.h):

字段含义
threadCount线程池中的线程总数
idleThreadCount空闲线程数
activeThreadCount活跃线程数
pendingTaskCount待执行任务数
totalTaskCount累计任务总数
processedTaskCount已处理任务数
maxIdleTime最长空闲时间

通过getPoolStats()获取;getPendingTaskCount()则单独返回待执行任务数。

如果需要任务级的观测(入队、出队、处理完成三个阶段的耗时),可以注册TaskObserver(接口见 folly/executors/ThreadPoolExecutor.h):

class TaskObserver { public: virtual ~TaskObserver() = default; virtual void taskEnqueued(const TaskInfo&) noexcept {} virtual void taskDequeued(const DequeuedTaskInfo&) noexcept {} virtual void taskProcessed(const ProcessedTaskInfo&) noexcept {} };

TaskInfo/DequeuedTaskInfo/ProcessedTaskInfo层层继承,分别携带优先级、requestId、入队时间、taskId,以及waitTime(出队时间 − 入队时间)和runTime(处理耗时)。任务处理完成后,runTask会回调所有已注册的TaskObserver::taskProcessed(见 folly/executors/ThreadPoolExecutor.cpp)。注意:TaskObserver出于性能考虑只能添加、不能移除,会在线程池析构时统一销毁;旧的subscribeToTaskStats(TaskStatsCallback)接口已被标记为 deprecated,建议迁移到addTaskObserver

选型建议与总结

回到文档的核心结论,选择线程池时可以遵循以下准则:

  1. IO 密集 / 需要事件循环(如 Thrift、memcache 客户端、异步 socket):使用IOThreadPoolExecutor,通过getEventBase()直接调度 IO 工作;默认每核一线程即可,除非线程会阻塞;
  2. CPU 密集计算(如 JSON 解析、加密、计算型回调):使用CPUThreadPoolExecutor,借助 LifoSem 的 LIFO 唤醒获得 cache locality,并可选用多优先级队列区分任务重要程度;
  3. 通用异步代码:直接使用全局 Executor(getGlobalCPUExecutor()/getGlobalIOExecutor())配合via()/then()组成 Future 流水线,避免每次手工创建线程池;
  4. 需要监控:通过getPoolStats()观察池级指标,通过addTaskObserver获取任务级 wait/run 耗时,通过 Observer 管理每线程资源;
  5. 注意停止语义差异IOThreadPoolExecutor::stop()类似join(),会等待事件循环中的任务完成;CPUThreadPoolExecutor::stop()则尽力完成未完成任务后返回;
  6. 避免有界阻塞队列引发死锁:优先使用无界队列,或只从线程池外部提交任务。

这套"IO 与 CPU 分离、事件循环与任务队列各司其职"的设计,是 Folly 得以在 Facebook 大规模生产环境中提供稳定尾延迟的基石之一。理解其背后的系统原语约束(event_fd、epoll、信号量、惊群),比记住 API 本身更能帮助你做出正确的并发架构决策。相关源码与测试(如 folly/executors/test/IOThreadPoolExecutorTest.cpp、folly/executors/test/ThreadPoolExecutorTest.cpp、folly/executors/test/GlobalExecutorTest.cpp)可作为进一步研究的入口。

【免费下载链接】follyAn open-source C++ library developed and used at Facebook.项目地址: https://gitcode.com/GitHub_Trending/fol/folly

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/10 11:37:07

TVBoxOSC 电视盒子控制使用完全指南:从安装到流畅播放的完整路线

TVBoxOSC 电视盒子控制使用完全指南&#xff1a;从安装到流畅播放的完整路线 【免费下载链接】TVBoxOSC TVBoxOSC - 一个基于第三方项目的代码库&#xff0c;用于电视盒子的控制和管理。 项目地址: https://gitcode.com/GitHub_Trending/tv/TVBoxOSC 晚上想看点片&#…

作者头像 李华
网站建设 2026/9/10 11:36:39

Python自行车共享需求预测实战:从数据清洗到可解释回归

简介&#xff1a;本资源是一份面向数据科学初学者与计算机相关专业学生的Kaggle实战项目&#xff0c;聚焦城市自行车共享系统使用状况的探索性分析与需求预测&#xff0c;适用于毕业设计、课程设计及算法入门实践。压缩包共8个文件&#xff0c;含3个核心数据集&#xff08;CSV&…

作者头像 李华