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

如何基于Redis++通过Pipeline批量获取Redis更新以提升性能?

解决Redis++订阅事件批量积累处理的问题

你遇到的核心问题是subscriber.consume()默认是阻塞式处理单个事件,每次调用只会处理一条到达的消息,导致无法批量积累键来提升pipeline的效率。下面提供两种可行的解决方案:

方案一:带超时的批量消费(单线程)

利用Redis++中subscriber.consume()的超时重载版本,在指定时间窗口内尽可能多地接收事件,积累足够数量的键后再执行批量查询。

修改后的代码示例

#include <sw/redis++/redis++.h>
#include <chrono>

using namespace sw::redis;

// 假设你的Entry类型和do_something_with函数已定义
struct Entry {
    std::string key;
    std::string value;
    std::size_t ttl;
    Entry(std::string k, std::string v, std::size_t t) : key(std::move(k)), value(std::move(v)), ttl(t) {}
};

void do_something_with(const Entry&);

int main() {
    Redis redis("tcp://127.0.0.1:6379");
    auto subscriber = redis.subscriber();
    subscriber.subscribe("__keyevent@0__:set");

    Pipeline pipeline = redis.pipeline();
    std::vector<std::string> updated_redis_keys;

    // 设置消息回调,积累更新的键
    subscriber.on_message([&updated_redis_keys](std::string, std::string key) {
        updated_redis_keys.push_back(std::move(key));
    });

    while (true) {
        try {
            // 等待100毫秒,期间处理所有到达的set事件
            subscriber.consume(std::chrono::milliseconds(100));

            if (updated_redis_keys.empty()) continue;

            // 批量构建pipeline命令
            for (const auto& key : updated_redis_keys) {
                pipeline.get(key).ttl(key);
            }

            auto replies = pipeline.exec();

            // 批量处理查询结果
            for (std::size_t i = 0; i < updated_redis_keys.size(); ++i) {
                static constexpr std::size_t ValueIndex = 0;
                static constexpr std::size_t TtlIndex = 1;

                const auto value = replies.get<std::optional<std::string>>(i * 2 + ValueIndex);
                const auto ttl = replies.get<long long>(i * 2 + TtlIndex);

                // 注意:必须判断value是否存在,避免key已被删除导致空指针
                if (value) {
                    Entry entry(std::move(updated_redis_keys[i]), *value, static_cast<std::size_t>(ttl));
                    do_something_with(entry);
                }
            }

            updated_redis_keys.clear();
            pipeline.clear(); // 清空pipeline,准备下一批命令
        } catch (const sw::redis::Error& err) {
            // 不要空捕获,至少输出错误信息便于调试
            std::cerr << "Redis操作错误: " << err.what() << std::endl;
        }
    }

    return 0;
}

关键说明

  • 超时时间可以根据业务场景调整:如果事件频率高,设为10-50ms;如果事件稀疏,可设为500ms-1s,平衡延迟和批量处理效率。
  • 每次处理完后要调用pipeline.clear(),避免残留上一批的命令。
  • 必须检查std::optional<std::string>是否有值,因为事件触发后key可能已被TTL或其他操作删除。

方案二:线程解耦(接收与处理分离)

将订阅事件的逻辑放到单独线程中,把更新的键存入线程安全队列,主线程定期从队列中批量取键执行查询。这种方式适合事件量较大的场景,避免查询阻塞事件接收。

实现步骤

  1. 定义线程安全队列
#include <queue>
#include <mutex>
#include <condition_variable>
#include <vector>
#include <utility>

template<typename T>
class ThreadSafeQueue {
public:
    void push(T item) {
        std::lock_guard<std::mutex> lock(_mtx);
        _queue.push(std::move(item));
        _cv.notify_one();
    }

    std::vector<T> pop_all() {
        std::lock_guard<std::mutex> lock(_mtx);
        std::vector<T> items;
        while (!_queue.empty()) {
            items.push_back(std::move(_queue.front()));
            _queue.pop();
        }
        return items;
    }

    bool empty() const {
        std::lock_guard<std::mutex> lock(_mtx);
        return _queue.empty();
    }

private:
    std::queue<T> _queue;
    mutable std::mutex _mtx;
    std::condition_variable _cv;
};
  1. 主线程与订阅线程分离的代码
#include <sw/redis++/redis++.h>
#include <thread>
#include <chrono>

using namespace sw::redis;

struct Entry {
    std::string key;
    std::string value;
    std::size_t ttl;
    Entry(std::string k, std::string v, std::size_t t) : key(std::move(k)), value(std::move(v)), ttl(t) {}
};

void do_something_with(const Entry&);

int main() {
    Redis redis("tcp://127.0.0.1:6379");
    auto subscriber = redis.subscriber();
    subscriber.subscribe("__keyevent@0__:set");

    Pipeline pipeline = redis.pipeline();
    ThreadSafeQueue<std::string> key_queue;

    // 启动订阅线程,专门接收事件
    std::thread sub_thread([&subscriber, &key_queue]() {
        subscriber.on_message([&key_queue](std::string, std::string key) {
            key_queue.push(std::move(key));
        });

        while (true) {
            try {
                // 阻塞接收消息,不影响主线程处理
                subscriber.consume();
            } catch (const sw::redis::Error& err) {
                std::cerr << "订阅线程错误: " << err.what() << std::endl;
                // 可选:添加重连逻辑,避免断开后无法接收事件
                std::this_thread::sleep_for(std::chrono::seconds(1));
                subscriber.subscribe("__keyevent@0__:set");
            }
        }
    });

    // 主线程批量处理键
    while (true) {
        // 每100ms检查一次队列,可根据需求调整间隔
        std::this_thread::sleep_for(std::chrono::milliseconds(100));

        auto keys = key_queue.pop_all();
        if (keys.empty()) continue;

        // 构建批量查询命令
        for (const auto& key : keys) {
            pipeline.get(key).ttl(key);
        }

        auto replies = pipeline.exec();

        // 处理结果
        for (std::size_t i = 0; i < keys.size(); ++i) {
            static constexpr std::size_t ValueIndex = 0;
            static constexpr std::size_t TtlIndex = 1;

            const auto value = replies.get<std::optional<std::string>>(i * 2 + ValueIndex);
            const auto ttl = replies.get<long long>(i * 2 + TtlIndex);

            if (value) {
                Entry entry(std::move(keys[i]), *value, static_cast<std::size_t>(ttl));
                do_something_with(entry);
            }
        }

        pipeline.clear();
    }

    // 程序结束时记得join线程(实际场景中需处理优雅退出)
    // sub_thread.join();

    return 0;
}

关键说明

  • 线程安全队列保证了订阅线程和主线程之间的键传递不会出现竞态条件。
  • 订阅线程如果断开连接,添加了重连逻辑,确保事件接收不中断。
  • 主线程的休眠间隔可根据业务延迟要求调整,间隔越短延迟越低,但批量大小可能越小。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 11:35:23