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
相关产品推荐
相关产品推荐

