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

C++17多线程问询:向量填充处理并行及API请求计算优化

C++17多线程并行优化方案:API请求与数据处理异步执行

需求与现有问题

  • 任务背景:需调用API获取数据(单请求耗时3秒),对返回数据执行耗时数十秒的本地计算,总请求数为M
  • 现有实现缺陷:创建N个线程,每个线程串行执行「API请求→数据解析→计算」流程,导致API等待时计算资源闲置,计算时请求资源闲置,完全未利用并行潜力
  • 优化目标:让API请求(含解析)与数据计算同时进行,采用LIFO(后进先出)结构存储待处理数据

核心技术方向与设计模式

  • 技术选型:基于C++17标准,使用std::thread、std::mutex、std::condition_variable实现线程同步,兼容Visual Studio 2019,无需第三方库
  • 设计模式:生产者-消费者模式(LIFO变体)
    • 生产者线程:专门负责发起API请求、解析数据,将结果压入线程安全的LIFO栈
    • 消费者线程:专门负责从LIFO栈中取出数据,执行耗时计算
    • LIFO选择:若后续请求的计算优先级更高(或无需顺序依赖),栈结构比队列更贴合需求

线程安全LIFO栈实现(核心通信组件)

#include <stack>
#include <mutex>
#include <condition_variable>
#include <optional>

// 线程安全的LIFO栈,用于存储待处理的JSON数据
template<typename T>
class ThreadSafeStack {
private:
    std::stack<T> data_stack;
    mutable std::mutex mtx;
    std::condition_variable cv;
    bool is_closed = false; // 标记是否停止接受新数据

public:
    // 压入数据,通知消费者线程
    void push(T value) {
        std::lock_guard<std::mutex> lock(mtx);
        data_stack.push(std::move(value));
        cv.notify_one();
    }

    // 弹出数据,无数据时阻塞等待,直到有数据或栈关闭
    std::optional<T> pop() {
        std::unique_lock<std::mutex> lock(mtx);
        cv.wait(lock, [this] { return is_closed || !data_stack.empty(); });
        
        if (is_closed && data_stack.empty()) {
            return std::nullopt;
        }

        T value = std::move(data_stack.top());
        data_stack.pop();
        return value;
    }

    // 关闭栈,通知所有消费者停止等待
    void close() {
        std::lock_guard<std::mutex> lock(mtx);
        is_closed = true;
        cv.notify_all();
    }
};

完整优化代码示例

#include <iostream>
#include <vector>
#include <thread>
#include <string>
#include <chrono>
// 引入上述ThreadSafeStack类

// 模拟API请求与解析,耗时3秒
std::string MakeRequestAndParse(int page) {
    std::this_thread::sleep_for(std::chrono::seconds(3));
    return "{\"page\":" + std::to_string(page) + ", \"data\": \"sample_content\"}";
}

// 模拟耗时计算,示例耗时10秒
void Calc(const std::string& json_data) {
    std::this_thread::sleep_for(std::chrono::seconds(10));
    std::cout << "已处理数据: " << json_data << std::endl;
}

// 生产者线程函数:处理指定范围的API请求
void Producer(ThreadSafeStack<std::string>& stack, int page_from, int page_to) {
    for (int i = page_from; i < page_to; ++i) {
        std::string json = MakeRequestAndParse(i);
        stack.push(std::move(json));
    }
}

// 消费者线程函数:从栈取数据并执行计算
void Consumer(ThreadSafeStack<std::string>& stack) {
    while (auto opt_data = stack.pop()) {
        Calc(*opt_data);
    }
}

int main() {
    const int PagesCount = 10; // 总请求数M
    const int ProducerThreads = 2; // 生产者线程数(根据API并发限制调整)
    const int ConsumerThreads = 3; // 消费者线程数(根据CPU核心数调整)
    const int BlockSize = (PagesCount + ProducerThreads - 1) / ProducerThreads; // 每个生产者处理的页数

    ThreadSafeStack<std::string> data_stack;
    std::vector<std::thread> producers;
    std::vector<std::thread> consumers;

    // 启动生产者线程
    int page_begin = 0;
    for (int i = 0; i < ProducerThreads; ++i) {
        int page_end = std::min(page_begin + BlockSize, PagesCount);
        producers.emplace_back(Producer, std::ref(data_stack), page_begin, page_end);
        page_begin = page_end;
    }

    // 启动消费者线程
    for (int i = 0; i < ConsumerThreads; ++i) {
        consumers.emplace_back(Consumer, std::ref(data_stack));
    }

    // 等待所有生产者完成请求,关闭栈
    for (auto& t : producers) {
        t.join();
    }
    data_stack.close();

    // 等待所有消费者完成计算
    for (auto& t : consumers) {
        t.join();
    }

    std::cout << "所有任务执行完成。" << std::endl;
    return 0;
}

关键优化点说明

  • 任务拆分:将串行的「请求→计算」拆分为独立的生产者(IO密集型)和消费者(CPU密集型)线程,实现两类任务并行执行
  • 线程安全通信:通过带条件变量的线程安全栈,避免忙等待,保证生产者和消费者的同步
  • 资源适配:生产者线程数可根据API并发限制调整,消费者线程数可根据CPU核心数灵活配置
  • LIFO特性:栈结构保证最新请求的数据优先被处理,符合需求要求

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 00:23:21