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

多线程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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 17:43:09