You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何获取boost::asio::thread_pool的任务队列大小?

Boost.Asio thread_pool 获取待执行任务队列大小方案

Boost.Asio的thread_pool本身并没有直接暴露获取内部任务队列大小的接口,想要实现这个需求,得自己做一层包装或者自定义线程池,下面给两种实用方案:

方案一:包装任务统计待完成任务数(简单易用)

这个方案通过原子变量统计已提交但未完成的任务总数(包括正在线程中执行的和队列里等待的),适合大多数流量控制场景。

封装带计数的线程池类

#include <boost/asio/thread_pool.hpp>
#include <atomic>
#include <functional>
#include <iostream>

class CountedThreadPool {
public:
    explicit CountedThreadPool(std::size_t num_threads) : pool_(num_threads) {}

    // 提交任务,自动更新计数
    template <typename Func>
    void post(Func&& func) {
        pending_tasks_.fetch_add(1, std::memory_order_relaxed);
        boost::asio::post(pool_, [this, func = std::forward<Func>(func)]() {
            try {
                func();
            } finally {
                // 无论任务是否抛出异常,都要减少计数
                pending_tasks_.fetch_sub(1, std::memory_order_relaxed);
            }
        });
    }

    // 获取当前待完成的任务数(含正在执行的)
    std::size_t pending_task_count() const {
        return pending_tasks_.load(std::memory_order_relaxed);
    }

    // 等待所有任务完成并销毁线程池
    void join() {
        pool_.join();
    }

private:
    boost::asio::thread_pool pool_;
    std::atomic<std::size_t> pending_tasks_{0};
};

修改业务代码使用该类

// 假设my_task是已定义的任务函数
void my_task() {
    // 任务逻辑...
}

// 异步回调函数,传递CountedThreadPool引用
static void somecallback(CountedThreadPool& pool) {
    const std::size_t kQueueLimit = 100; // 自定义队列上限
    auto current_count = pool.pending_task_count();
    if (current_count < kQueueLimit) {
        pool.post(my_task);
        // 日志输出当前计数
        std::cout << "提交任务成功,当前待完成任务数:" << pool.pending_task_count() << std::endl;
    } else {
        std::cerr << "队列溢出!当前待完成任务数:" << current_count << std::endl;
    }
}

// 主线程中的类
class A {
public:
    A() : pool_(4) {};
    ~A() {
        pool_.join();
    }
    // 其他业务逻辑...
private:
    CountedThreadPool pool_;
};

方案二:自定义线程池获取严格队列待执行数(精确需求)

如果需要严格区分队列中等待执行的任务数(不包含正在执行的),就得自己实现线程池的任务调度逻辑,因为boost的thread_pool内部队列是私有的,无法直接访问。

自定义线程池实现

#include <queue>
#include <mutex>
#include <condition_variable>
#include <vector>
#include <thread>
#include <functional>
#include <iostream>

class CustomThreadPool {
public:
    explicit CustomThreadPool(std::size_t num_threads) {
        // 启动指定数量的工作线程
        for (std::size_t i = 0; i < num_threads; ++i) {
            threads_.emplace_back([this]() {
                while (true) {
                    std::function<void()> task;
                    {
                        std::unique_lock<std::mutex> lock(mutex_);
                        // 等待任务或停止信号
                        cv_.wait(lock, [this]() { return !tasks_.empty() || stopped_; });
                        if (stopped_ && tasks_.empty()) break;
                        // 取出队列中的任务
                        task = std::move(tasks_.front());
                        tasks_.pop();
                    }
                    // 执行任务
                    task();
                }
            });
        }
    }

    ~CustomThreadPool() {
        {
            std::lock_guard<std::mutex> lock(mutex_);
            stopped_ = true;
        }
        cv_.notify_all();
        // 等待所有线程退出
        for (auto& t : threads_) {
            t.join();
        }
    }

    // 提交任务到队列
    template <typename Func>
    void post(Func&& func) {
        {
            std::lock_guard<std::mutex> lock(mutex_);
            tasks_.emplace(std::forward<Func>(func));
        }
        cv_.notify_one();
    }

    // 获取队列中待执行的任务数(不包含正在执行的)
    std::size_t queue_size() const {
        std::lock_guard<std::mutex> lock(mutex_);
        return tasks_.size();
    }

private:
    mutable std::mutex mutex_;
    std::condition_variable cv_;
    std::queue<std::function<void()>> tasks_;
    std::vector<std::thread> threads_;
    bool stopped_ = false;
};

使用示例

void my_task() {
    // 任务逻辑...
}

static void somecallback(CustomThreadPool& pool) {
    const std::size_t kQueueLimit = 100;
    auto current_queue_size = pool.queue_size();
    if (current_queue_size < kQueueLimit) {
        pool.post(my_task);
        std::cout << "提交任务成功,当前队列待执行任务数:" << pool.queue_size() << std::endl;
    } else {
        std::cerr << "队列溢出!当前队列待执行任务数:" << current_queue_size << std::endl;
    }
}

class A {
public:
    A() : pool_(4) {};
    ~A() {
        // 自定义线程池的析构函数会自动处理join
    }
    // 其他业务逻辑...
private:
    CustomThreadPool pool_;
};

注意事项

  • 方案一的计数包含正在执行的任务,实现简单,适合不需要严格区分队列等待数的场景;
  • 方案二可以精确获取队列中待执行的任务数,但需要自己维护线程和队列,代码量更大;
  • 两种方案的计数都可以直接用于日志输出,方便监控队列状态。

内容的提问来源于stack exchange,提问作者Mikhail

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 17:47:03