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

如何在Cap'n Proto中实现类似Boost的阻塞重试或异步回调重连功能?

刚好我之前在Cap'n Proto里实现过类似的连接重试逻辑,不管是阻塞式的(像Boost的connect那样一直重试到成功)还是异步回调式的(类似basic_socket::async_connect)都能搞定,给你详细说说:

阻塞式重试连接(仿Boost风格)

如果你想要一个阻塞直到连接成功(或达到最大重试次数)的函数,可以通过循环捕获连接异常,结合kj::sleep实现间隔重试。这里直接封装成一个复用函数:

#include <capnp/ez-rpc.h>
#include <kj/async.h>
#include <iostream>
#include <chrono>

kj::Own<kj::AsyncIoStream> connectWithRetry(
    kj::AsyncIoContext& ioContext,
    kj::StringPtr address,
    unsigned int maxRetries,  // 0表示无限重试
    std::chrono::milliseconds retryDelay
) {
    kj::WaitScope& waitScope = ioContext.waitScope;
    auto network = ioContext.provider->getNetwork();
    
    for (unsigned int attempt = 0; maxRetries == 0 || attempt < maxRetries; ++attempt) {
        try {
            // 解析地址并尝试连接
            auto parsedAddr = network.parseAddress(address).wait(waitScope);
            return parsedAddr->connect().wait(waitScope);
        } catch (const kj::Exception& e) {
            std::cerr << "连接尝试 #" << (attempt + 1) << " 失败: " << e.getDescription() << std::endl;
            
            // 如果是最后一次重试,直接抛出异常
            if (maxRetries != 0 && attempt == maxRetries - 1) {
                throw;
            }
            
            // 等待指定间隔后重试
            kj::sleep(retryDelay).wait(waitScope);
        }
    }
    
    KJ_UNREACHABLE;  // 无限重试时不会走到这里
}

// 使用示例
int main() {
    auto ioContext = kj::setupAsyncIo();
    
    try {
        // 无限重试直到连接成功(也可以指定比如maxRetries=10)
        auto connection = connectWithRetry(ioContext, "localhost:7500", 0, std::chrono::milliseconds(1000));
        std::cout << "连接成功!" << std::endl;
        
        // 这里可以继续处理连接后的逻辑,比如创建RPC客户端等
    } catch (const kj::Exception& e) {
        std::cerr << "所有连接尝试均失败: " << e.getDescription() << std::endl;
        return 1;
    }
    
    return 0;
}

关键点说明:

  • maxRetries设为0时会无限重试,直到连接成功;指定具体数值则在达到次数后抛出异常。
  • 捕获kj::Exception:Cap'n Proto的异步操作抛出的异常都继承自这个类,能覆盖所有连接相关的错误(比如端口未监听、网络不可达等)。
  • kj::sleep用来控制重试间隔,避免频繁尝试占用资源。
异步回调式重连(仿async_connect风格)

如果想要非阻塞的异步重连,利用Cap'n Proto的Promise链式调用就能实现,通过回调通知连接成功或最终失败:

#include <capnp/ez-rpc.h>
#include <kj/async.h>
#include <iostream>
#include <chrono>

// 定义回调类型
using OnConnectSuccess = kj::Function<void(kj::Own<kj::AsyncIoStream>)>;
using OnConnectFailure = kj::Function<void(const kj::Exception&)>;

void asyncConnectWithRetry(
    kj::AsyncIoContext& ioContext,
    kj::StringPtr address,
    unsigned int maxRetries,
    std::chrono::milliseconds retryDelay,
    OnConnectSuccess onSuccess,
    OnConnectFailure onFailure
) {
    auto network = ioContext.provider->getNetwork();
    kj::WaitScope& waitScope = ioContext.waitScope;
    
    // 递归的连接尝试函数
    auto attemptConnect = [&](unsigned int attempt) -> kj::Promise<kj::Own<kj::AsyncIoStream>> {
        return network.parseAddress(address)
            .then([](kj::Own<kj::NetworkAddress> addr) {
                return addr->connect();
            })
            .catch_([&](kj::Exception&& e) -> kj::Promise<kj::Own<kj::AsyncIoStream>> {
                std::cerr << "异步连接尝试 #" << (attempt + 1) << " 失败: " << e.getDescription() << std::endl;
                
                // 重试次数耗尽,返回失败的Promise
                if (maxRetries != 0 && attempt >= maxRetries - 1) {
                    return kj::mv(e);
                }
                
                // 等待后递归重试
                return kj::sleep(retryDelay)
                    .then([&, attempt]() {
                        return attemptConnect(attempt + 1);
                    });
            });
    };
    
    // 启动连接尝试链,并绑定回调
    attemptConnect(0)
        .then([onSuccess = kj::mv(onSuccess)](kj::Own<kj::AsyncIoStream> conn) {
            onSuccess(kj::mv(conn));
        })
        .catch_([onFailure = kj::mv(onFailure)](kj::Exception&& e) {
            onFailure(e);
        })
        .attach(waitScope);  // 确保Promise在waitScope的生命周期内执行
}

// 使用示例
int main() {
    auto ioContext = kj::setupAsyncIo();
    
    asyncConnectWithRetry(
        ioContext,
        "localhost:7500",
        5,  // 最多重试5次
        std::chrono::milliseconds(1000),
        [](kj::Own<kj::AsyncIoStream> conn) {
            std::cout << "异步连接成功!" << std::endl;
            // 处理连接成功后的逻辑,比如初始化RPC客户端
        },
        [](const kj::Exception& e) {
            std::cerr << "所有异步连接尝试均失败: " << e.getDescription() << std::endl;
            // 处理最终失败的逻辑
        }
    );
    
    // 启动事件循环,等待异步操作完成
    ioContext.waitScope.waitForNoWork();
    
    return 0;
}

关键点说明:

  • 用递归的attemptConnect函数构建Promise链,失败后自动重试。
  • 通过then和catch_处理成功/失败的回调,完全非阻塞。
  • attach(waitScope)是必须的,它能保证Promise不会被提前销毁,直到异步操作完成。
  • 如果需要支持取消重试,可以结合kj::CancellationSet,给sleep和connect操作添加取消令牌,这样就能在外部触发取消逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:50:16