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

自定义boost::asio异步操作实现及相关技术问题咨询

自定义Boost.Asio异步操作相关问题

我正尝试实现自己的boost::asio异步操作,使其在io_context.run()期间异步执行。同时需要实现一个可取消的操作,该操作会等待某个条件/谓词满足(或boost::signals2::signal被触发),之后在io_context线程中调用完成处理器。想知道Boost中是否已提供此类功能?

我参考Boost文档写了以下代码:

#include <boost/asio.hpp>
#include <iostream>

/// From boost documentation:
/// — If the initiating function is not a member function,
/// the associated executor is that returned by the get_executor member function
/// of the first argument to the initiating function.
struct executor_owner
{
    using executor_type = boost::asio::io_context;
    executor_type* m_executor = nullptr;

    explicit executor_owner(executor_type& ex)
        : m_executor(&ex)
    {
    }

    executor_type& get_executor() BOOST_ASIO_NOEXCEPT
    {
        std::cout << "get_executor() called." << std::endl;
        return *m_executor;
    }
};

/// My initiating function
template<class CompletionToken>
auto async_xyz(executor_owner executorOwner, CompletionToken&& token)
{
    using completion_handler_t = typename boost::asio::async_completion<CompletionToken, void(executor_owner&)>::completion_handler_type;
    return boost::asio::async_initiate<CompletionToken, void(executor_owner&)>(
        [](completion_handler_t completion_handler, executor_owner& owner) {
            // initiate the operation and cause completion_handler to be invoked
            // with the result
            std::cout << "I'm performing the async func operation here!" << std::endl;
            completion_handler();
        },
        token, executorOwner);
}

int main()
{
    boost::asio::io_context io_context;
    executor_owner executorOwner(io_context);

    async_xyz(executorOwner, []()
    {
        std::cout << "completion handler here!" << std::endl;
    });

    // I expect that completion handler will be called asynchronously, during io_context::run().
    // Unfortunately, completion handler has already been called.

    // std::cout << "io_context begin." << std::endl;
    // io_context.run(); // It doesn't matter. completion handler
    // std::cout << "io_context done." << std::endl;

    return 0;
}

输出结果:

I'm performing the async func operation here!
completion handler here!

我的问题:

  1. 如何实现符合boost::asio设计规范的自定义异步操作?
  2. 为什么我的executor没有按照Boost文档所述通过get_executor()方法获取?
  3. boost::asio::async_initiate内部是否会将操作提交到io_context,还是需要我在自定义操作实现中手动调用io_context::post(...)来触发完成处理器?
  4. 如何实现一个简单的等待操作,用于等待谓词满足或信号触发?
  5. 如何为此类操作添加超时/取消功能?

使用环境:Boost 1.82、C14(优先考虑C11解决方案)


问题解答

1. 符合Boost.Asio规范的自定义异步操作实现方式

需遵循组合异步操作的设计模式,核心要点:

  • 利用boost::asio::async_completion和async_initiate处理CompletionToken(支持回调、future等多种形式)
  • 确保操作异步性:必须将完成处理器提交到关联的executor(如io_context)执行,禁止直接同步调用
  • 遵循Executor模型:通过get_associated_executor获取与完成处理器关联的executor,或从操作第一个参数(如你的executor_owner)获取
  • 保证异常安全:操作过程中若发生异常,需正确传递给完成处理器,避免资源泄漏

2. get_executor()未被调用的原因

你的代码中async_initiate的lambda直接同步调用了completion_handler,且未触发Executor关联逻辑:

  • async_initiate本身不会自动获取executor,需要手动在lambda中通过boost::asio::get_associated_executor获取,或显式调用executor_owner的get_executor()
  • 你没有利用executor_owner将任务提交到io_context,因此get_executor()没有触发场景

修正示例:

auto ex = boost::asio::get_associated_executor(completion_handler, owner.get_executor());
ex.post([completion_handler]() mutable {
    completion_handler();
});

3. async_initiate是否自动提交到io_context?

不会。async_initiate仅负责处理CompletionToken的类型擦除,将用户提供的token转换为统一的完成处理器类型,不会自动将操作或完成处理器提交到executor。必须手动调用executor的post()/dispatch()/defer()方法,将完成处理器的执行提交到io_context线程,才能实现异步效果。

你的代码直接同步调用completion_handler(),因此会在调用async_xyz的线程立即执行,而非等待io_context.run()。

4. 实现等待谓词/信号触发的异步操作

方式1:等待谓词满足

使用定时器定期检查谓词,条件满足后触发完成处理器:

template <typename Predicate, typename CompletionToken>
auto async_wait_predicate(boost::asio::io_context& io, Predicate pred, CompletionToken&& token)
{
    using completion_handler_t = typename boost::asio::async_completion<CompletionToken, void()>::completion_handler_type;
    
    return boost::asio::async_initiate<CompletionToken, void()>(
        [&io, pred](completion_handler_t handler) mutable {
            auto check_predicate = [&io, pred, handler](const boost::system::error_code& ec) mutable {
                if (!ec && pred()) {
                    boost::asio::post(io, std::move(handler));
                    return;
                }
                // 10ms后再次检查
                auto timer = std::make_shared<boost::asio::steady_timer>(io);
                timer->expires_after(std::chrono::milliseconds(10));
                timer->async_wait(std::move(check_predicate));
            };
            check_predicate({});
        },
        token
    );
}

方式2:等待boost::signals2::signal触发

将完成处理器绑定到信号,触发时提交到io_context执行:

template <typename CompletionToken>
auto async_wait_signal(boost::asio::io_context& io, boost::signals2::signal<void()>& sig, CompletionToken&& token)
{
    using completion_handler_t = typename boost::asio::async_completion<CompletionToken, void()>::completion_handler_type;
    
    return boost::asio::async_initiate<CompletionToken, void()>(
        [&io, &sig](completion_handler_t handler) {
            sig.connect([&io, handler]() mutable {
                boost::asio::post(io, std::move(handler));
            });
        },
        token
    );
}

5. 添加超时/取消功能

超时功能

在等待操作中添加定时器,超时后触发取消逻辑:

template <typename Predicate, typename CompletionToken>
auto async_wait_predicate_with_timeout(boost::asio::io_context& io, Predicate pred, std::chrono::milliseconds timeout, CompletionToken&& token)
{
    using completion_handler_t = typename boost::asio::async_completion<CompletionToken, void(boost::system::error_code)>::completion_handler_type;
    
    return boost::asio::async_initiate<CompletionToken, void(boost::system::error_code)>(
        [&io, pred, timeout](completion_handler_t handler) mutable {
            auto timer = std::make_shared<boost::asio::steady_timer>(io);
            auto timeout_timer = std::make_shared<boost::asio::steady_timer>(io, timeout);
            
            auto check_predicate = [&io, pred, handler, timer, timeout_timer](const boost::system::error_code& ec) mutable {
                if (ec == boost::asio::error::operation_aborted) {
                    handler(boost::asio::error::operation_aborted);
                    return;
                }
                if (pred()) {
                    timeout_timer->cancel();
                    boost::asio::post(io, [handler]() { handler({}); });
                    return;
                }
                timer->expires_after(std::chrono::milliseconds(10));
                timer->async_wait(std::move(check_predicate));
            };
            
            timeout_timer->async_wait([handler, timer](const boost::system::error_code& ec) {
                if (!ec) {
                    timer->cancel();
                    handler(boost::asio::error::timed_out);
                }
            });
            
            check_predicate({});
        },
        token
    );
}

取消功能

使用boost::asio::cancellation_signal,触发信号时终止等待操作:

template <typename Predicate, typename CompletionToken>
auto async_wait_predicate_cancellable(boost::asio::io_context& io, Predicate pred, boost::asio::cancellation_signal& cancel_sig, CompletionToken&& token)
{
    using completion_handler_t = typename boost::asio::async_completion<CompletionToken, void(boost::system::error_code)>::completion_handler_type;
    
    return boost::asio::async_initiate<CompletionToken, void(boost::system::error_code)>(
        [&io, pred, &cancel_sig](completion_handler_t handler) mutable {
            auto timer = std::make_shared<boost::asio::steady_timer>(io);
            auto slot = cancel_sig.connect([timer]() { timer->cancel(); });
            
            auto check_predicate = [&io, pred, handler, timer, slot](const boost::system::error_code& ec) mutable {
                if (ec == boost::asio::error::operation_aborted) {
                    handler(boost::asio::error::operation_aborted);
                    return;
                }
                if (pred()) {
                    boost::asio::post(io, [handler]() { handler({}); });
                    return;
                }
                timer->expires_after(std::chrono::milliseconds(10));
                timer->async_wait(std::move(check_predicate));
            };
            
            check_predicate({});
        },
        token
    );
}

内容的提问来源于stack exchange,提问作者Michał Jaroń

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 22:15:53