多线程RPC请求合并批量调用及结果分发的C++实现思路与代码咨询
批量RPC请求合并与结果分发实现方案
核心思路
要实现多线程请求的批量合并与结果分发,关键在于以下几点:
- 请求收集与同步:使用线程安全的队列存储每个线程的请求,通过互斥锁和条件变量保证队列操作的原子性,唤醒批量处理线程。
- 请求关联追踪:每个请求绑定一个
std::promise,线程通过对应的std::future等待结果;批量处理时记录每个请求的输入长度,以便从RPC响应中截取对应结果。 - 批量处理线程:定期从队列中取出所有待处理请求,合并为单个RPC调用,将响应结果拆分后通过
promise返回给对应线程。
C++示例代码
1. 基础结构定义
首先定义与RPC IDL匹配的请求、响应结构,并模拟RPC调用:
#include <vector> #include <thread> #include <mutex> #include <condition_variable> #include <future> #include <queue> #include <iostream> #include <chrono> // 匹配RPC IDL的请求结构 struct Req { std::vector<int> input; }; // 匹配RPC IDL的响应结构 struct Rsp { std::vector<int> output; }; // 模拟批量RPC调用(实际场景替换为真实RPC客户端调用) Rsp batch_rpc_call(const Req& req) { Rsp rsp; for (int num : req.input) { rsp.output.push_back(num * num); } // 模拟网络延迟 std::this_thread::sleep_for(std::chrono::milliseconds(100)); return rsp; }
2. 批量RPC管理器
实现单例模式的管理器,负责请求收集、批量处理和结果分发:
class BatchRpcManager { private: // 队列元素:输入列表 + 用于返回结果的promise using RequestEntry = std::pair<std::vector<int>, std::promise<std::vector<int>>>; std::queue<RequestEntry> request_queue; std::mutex queue_mutex; std::condition_variable cv; std::thread batcher_thread; bool running = true; // 私有构造函数,启动批量处理线程 BatchRpcManager() { batcher_thread = std::thread(&BatchRpcManager::batcher_loop, this); } // 批量处理线程主循环 void batcher_loop() { while (running) { std::vector<RequestEntry> collected_requests; // 加锁收集所有待处理请求 { std::unique_lock<std::mutex> lock(queue_mutex); // 等待队列非空或管理器停止 cv.wait(lock, [this]() { return !request_queue.empty() || !running; }); if (!running && request_queue.empty()) break; // 取出队列中所有请求 while (!request_queue.empty()) { collected_requests.push_back(std::move(request_queue.front())); request_queue.pop(); } } if (collected_requests.empty()) continue; // 构建批量请求 Req batch_req; std::vector<size_t> request_lengths; for (const auto& entry : collected_requests) { request_lengths.push_back(entry.first.size()); batch_req.input.insert(batch_req.input.end(), entry.first.begin(), entry.first.end()); } // 调用批量RPC Rsp batch_rsp = batch_rpc_call(batch_req); // 分发结果到每个请求的promise size_t current_pos = 0; for (size_t i = 0; i < collected_requests.size(); ++i) { size_t len = request_lengths[i]; std::vector<int> result(batch_rsp.output.begin() + current_pos, batch_rsp.output.begin() + current_pos + len); collected_requests[i].second.set_value(result); current_pos += len; } } } public: // 获取单例实例 static BatchRpcManager& get_instance() { static BatchRpcManager instance; return instance; } // 禁止拷贝和赋值 BatchRpcManager(const BatchRpcManager&) = delete; BatchRpcManager& operator=(const BatchRpcManager&) = delete; // 提交请求,返回用于等待结果的future std::future<std::vector<int>> submit_request(const std::vector<int>& input) { std::promise<std::vector<int>> prom; std::future<std::vector<int>> fut = prom.get_future(); { std::lock_guard<std::mutex> lock(queue_mutex); request_queue.emplace(input, std::move(prom)); } cv.notify_one(); // 唤醒批量处理线程 return fut; } // 析构函数:停止批量线程 ~BatchRpcManager() { running = false; cv.notify_one(); if (batcher_thread.joinable()) { batcher_thread.join(); } } };
3. 线程客户端示例
编写三个测试线程,模拟请求提交与结果等待:
// 线程执行函数 void worker_thread(int thread_id, const std::vector<int>& input) { auto future_result = BatchRpcManager::get_instance().submit_request(input); // 阻塞等待结果 std::vector<int> result = future_result.get(); // 输出结果 std::cout << "线程" << thread_id << "收到结果:"; for (int num : result) { std::cout << num << " "; } std::cout << std::endl; } int main() { // 创建三个测试线程 std::thread t1(worker_thread, 1, {1}); std::thread t2(worker_thread, 2, {2}); std::thread t3(worker_thread, 3, {3, 4}); // 等待所有线程完成 t1.join(); t2.join(); t3.join(); return 0; }
代码说明
- 单例管理器:确保全局只有一个批量处理线程,避免资源浪费。
- Promise/Future机制:每个线程提交请求后获得
future,通过get()阻塞等待结果;批量线程完成RPC调用后,通过promise.set_value()将结果返回给对应线程。 - 线程安全队列:使用
std::mutex和std::condition_variable保证队列操作的原子性,批量线程仅在有请求时被唤醒。
运行该代码,三个线程的请求会被合并为一次RPC调用,最终输出:
线程1收到结果:1 线程2收到结果:4 线程3收到结果:9 16
内容的提问来源于stack exchange,提问作者muller john
相关产品推荐
相关产品推荐

