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

