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

使用可等待运算符无法取消azmq::async_receive操作的问题

问题:Azmq async_receive 在定时器触发后未被取消

注:不确定这是Asio Bug、Azmq Bug还是自身使用问题。

我正在使用Azmq库从socket读取数据,该库的自定义socket类型azmq::socket是boost::asio::socket的封装。执行流程如下:

参考定时器函数

auto timer(std::chrono::steady_clock::duration dur) -> asio::awaitable<void> {
    asio::steady_timer timer(co_await asio::this_coro::executor);
    timer.expires_after(dur);
    co_await timer.async_wait(asio::use_awaitable);
}

合并等待的核心代码

auto val = co_await(azmq::async_receive(socket_, asio::buffer(buffer), asio::use_awaitable) || timer(std::chrono::seconds{1}));

Azmq async_receive 的相关实现

顶层async_receive模板

template<class CompletionToken, class MutableBufferSequence>
auto async_receive(azmq::socket &socket, MutableBufferSequence const &buffers, CompletionToken &&token) -> BOOST_ASIO_INITFN_RESULT_TYPE(CompletionToken, void(boost::system::error_code, size_t)) {
    return boost::asio::async_initiate<CompletionToken, void(boost::system::error_code, size_t)>(async_receive_initiation<MutableBufferSequence>{socket, buffers}, token);
}

异步操作初始化器

template<typename MutableBufferSequence>
struct async_receive_initiation {
  azmq::socket &socket;
  MutableBufferSequence const &buffers;

  template<typename CompletionHandler>
  void operator()(CompletionHandler &&completion_handler) {
      auto executor = boost::asio::get_associated_executor(
          completion_handler, socket.get_executor());
      socket.async_receive(buffers, boost::asio::bind_executor(executor,
                                                            std::bind(std::forward<CompletionHandler>(completion_handler), std::placeholders::_1, std::placeholders::_2)));
  }
};

底层socket的async_receive实现

template<typename MessageReadHandler>
void async_receive(MessageReadHandler && handler, flags_type flags = 0) {
        using type = detail::receive_op<MessageReadHandler>;
        get_service().enqueue<type>(get_implementation(), detail::socket_service::op_type::read_op, std::forward<MessageReadHandler>(handler), flags);
}

问题现象

定时器触发后,async_receive操作并未被取消。查阅Asio文档未找到解决办法,尝试修改完成令牌、给异步处理结构体添加取消函数均无效。

最小复现代码

#pragma once

#include <azmq/socket.hpp>
#include <boost/asio.hpp>
#include <boost/asio/experimental/awaitable_operators.hpp>
#include <iostream>

namespace asio = boost::asio;
using namespace asio::experimental::awaitable_operators;

auto timer(std::chrono::steady_clock::duration dur) -> asio::awaitable<void> {
    asio::steady_timer timer(co_await asio::this_coro::executor);
    timer.expires_after(dur);
    co_await timer.async_wait(asio::use_awaitable);
}

auto foo(asio::io_context &ctx) -> asio::awaitable<void> {

    azmq::sub_socket socket_{ctx};
    socket_ = azmq::sub_socket(socket_.get_io_context(), true);
    boost::system::error_code error_code;
    std::string const socket_path{"ipc:///tmp/not_real"};
    if (socket_.connect(socket_path, error_code)) {
        std::exit(-1);
    }
    if (socket_.set_option(azmq::socket::subscribe(""), error_code)) {
        std::exit(-1);
    }

    std::array<std::byte, 1024> buffer{};

    std::cout << "waiting for async_receive" << std::endl;
    co_await(azmq::async_receive(socket_, asio::buffer(buffer), asio::use_awaitable) || timer(std::chrono::seconds{1}));

    std::cout << "async_receive finished" << std::endl;

}

auto main() -> int {

    boost::asio::io_context ctx{};
    boost::asio::co_spawn(ctx, foo(ctx), asio::detached);
    ctx.run();

    return 0;
}

问题分析与解决思路

Asio的awaitable_operators::||操作符在其中一个awaitable完成时,会尝试取消另一个可取消的异步操作,但这个机制生效的前提是异步操作实现了Asio的取消模型:

  1. 异步操作需要关联取消令牌(Cancellation Token),并监听取消请求;
  2. 底层操作需要支持主动取消,或者在执行前检查取消状态。

从Azmq的代码实现来看,问题出在:

  • azmq::async_receive的初始化器没有获取或关联完成处理程序的取消槽(cancellation slot),无法接收Asio的取消信号;
  • 底层的socket服务只是将接收操作入队,没有提供取消已入队操作的接口,也不会在操作执行前检查是否已被取消。

可行解决方案

  1. 修改Azmq库以支持取消:
    在async_receive_initiation的operator()中,获取完成处理程序的关联取消槽,将其与socket的操作绑定。如果Azmq的socket服务支持取消已入队操作,就在取消触发时调用对应接口;否则,在操作执行前检查取消状态,若已取消则直接返回boost::asio::error::operation_aborted。

    示例修改方向:

    template<typename CompletionHandler>
    void operator()(CompletionHandler &&completion_handler) {
        auto executor = boost::asio::get_associated_executor(
            completion_handler, socket.get_executor());
        auto cancellation_slot = boost::asio::get_associated_cancellation_slot(
            completion_handler, socket.get_executor().context());
    
        // 封装处理程序,添加取消检查或绑定取消回调
        auto wrapped_handler = [completion_handler = std::forward<CompletionHandler>(completion_handler), cancellation_slot](boost::system::error_code ec, size_t size) mutable {
            if (cancellation_slot.is_cancelled()) {
                ec = boost::asio::error::operation_aborted;
            }
            std::move(completion_handler)(ec, size);
        };
    
        socket.async_receive(buffers, boost::asio::bind_executor(executor, std::move(wrapped_handler)));
    }
    

    注意:这只是简化示例,实际需要Azmq底层支持取消已入队的read_op才能完全生效。

  2. 手动封装支持取消的接收操作:
    若无法修改Azmq库,可以手动用asio::cancellation_signal实现取消逻辑:在定时器完成时发送取消信号,同时在接收操作的处理程序中检查该信号。不过这种方式需要手动管理信号,且无法真正取消Azmq底层的等待操作,只能提前终止处理流程。

  3. 使用Azmq原生超时选项(仅同步场景适用):
    Azmq的socket支持ZMQ_RCVTIMEO选项设置同步接收的超时,但这不适用于异步接收场景,无法直接解决当前问题。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 04:42:03