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

关于async_compose未维持io_context工作及Asio封装阻塞函数的问题

Asio异步封装阻塞函数的常见问题解答

我正在实现一个通用的async_task(Executor& executor, Token&& token, Fn&& func, Args&&... args)异步启动函数,目标是将任意第三方阻塞函数封装到线程中,并提供基于Asio的接口。目前实现尚未完善(例如已知需将完成处理程序提交到执行器而非在线程中运行),但已接近目标。我遇到三个问题:

  • 为什么程序会在所有完成处理程序运行前停止?原以为异步操作进行中无需工作守卫,但实际非回调处理程序完全未被调用,这是为什么?
  • 此代码是否违反Asio的设计原则?这类需求看似常见,但相关示例极少。
    -(附加问题)我想将std::thread替换为concurrency::task<void>,但无法在lambda捕获中使用仅移动类型。尝试self = make_shared<remove_reference_t<Self>>(move(self))后,三个处理程序仅输出str: 而无参数,推测与Self类型(实际为asio::detail::compose_op)包含已移入的impl副本有关,调用时使用了被移动后的旧版本,请问原因是什么?

代码实现

#include <chrono>
#include <iostream>
#include <memory>
#include <thread>

#include "asio.hpp"

template <typename Fn, typename... Args>
struct async_task_impl {
  std::decay_t<Fn> fn_;
  std::tuple<std::decay_t<Args>...> args_;

  async_task_impl(Fn&& fn, Args&&... args)
      : fn_(std::forward<Fn>(fn)), args_(std::forward<Args>(args)...) {}

  template <typename Self>
  auto operator()(Self& self) {
    // @todo: use concurrency::create_task
    auto t =
        std::thread([me = *this,             // copy impl into thread
                     self = std::move(self)  // move composed_op into thread?
    ]() mutable {
          try {
            std::apply(me.fn_, me.args_);
            self.complete({});
          } catch (std::exception& e) {
            self.complete(std::current_exception());
          }
        });
    t.detach();
  }
};

// runs some blocking task on its own thread and wraps it in asio
template <typename Executor, typename Token, typename Fn, typename... Args>
auto async_task(Executor& executor, Token&& token, Fn&& func, Args&&... args) {
  return asio::async_compose<Token, void(std::exception_ptr)>(
      async_task_impl(std::forward<Fn>(func), std::forward<Args>(args)...),
      token, executor);
}

测试代码

void slow_print(std::string str) {
  static std::mutex m;
  std::this_thread::sleep_for(std::chrono::milliseconds(500));
  {
    std::unique_lock lk(m);
    std::cout << "slow_print: " << str << "\n";
  }
  std::this_thread::sleep_for(std::chrono::milliseconds(500));
}

int main() {
  try {
    asio::io_context ctx;

    using namespace std::string_literals;

    async_task(
        ctx, [](std::exception_ptr) { std::cout << "callback done\n"; },
        slow_print, "callback"s);

    asio::co_spawn(
        ctx,
        [&]() -> asio::awaitable<void> {
          co_await async_task(ctx, asio::use_awaitable, slow_print, "coro"s);
        },
        asio::detached);

    auto f = std::async(std::launch::async, [&] {
      auto fut = async_task(ctx, asio::use_future, slow_print, "future"s);
      fut.get();
    });

    ctx.run();
  } catch (std::exception& e) {
    std::cout << e.what() << "\n";
  }
  return 0;
}

问题解答

1. 程序提前停止的原因

你的async_task_impl启动线程后直接返回,此时asio::async_compose会认为异步操作已完成,不会向io_context添加任何工作项。ctx.run()启动后,只会处理当前就绪的handler(比如async_task初始化逻辑),一旦没有更多工作就会退出,根本等不到后台线程调用self.complete()。

另外,use_awaitable和use_future对应的操作对象,在completion handler未被调度前,不会让io_context保持工作状态——只有当handler被提交到io_context时,才会关联工作。但你的后台线程调用self.complete()时,io_context已经停止,这些handler自然不会执行。

修复核心:必须让async_compose知道操作未完成,直到后台任务结束;同时要将self.complete()的调用提交到执行器,而非在后台线程直接调用(违反Asio线程安全规则)。

2. 是否违反Asio设计原则

需求本身是合理的,但当前实现存在两处违反Asio设计原则的问题:

  • Handler调度违规:直接在后台线程调用self.complete()绕过了Asio的执行器调度,违反了"handler必须在关联执行器上运行"的原则,可能导致io_context内部状态被并发访问,引发线程安全问题。
  • 工作管理缺失:没有正确维护io_context的工作状态,导致其提前退出。

Asio官方其实有类似封装思路(比如asio::post配合use_future/use_awaitable),正确的实现应该:

  1. 使用线程池而非临时创建线程,避免资源浪费。
  2. 任务完成后,通过执行器将completion handler调度到io_context。
  3. 确保async_compose的操作对象在任务完成前保持存活,并正确关联工作。

3. 替换为concurrency::task<void>后的参数丢失问题

你尝试将self放入shared_ptr时,复制的是已经被移动的asio::detail::compose_op对象。compose_op内部持有async_task_impl的副本,移动操作会转移原对象的内部资源,此时复制的是一个空/无效的compose_op,导致后续调用时参数丢失。

正确的做法是直接将self移动到shared_ptr中,而非复制已移动的对象,同时确保complete()在执行器上调用:

template <typename Self>
auto operator()(Self& self) {
  auto self_ptr = std::make_shared<Self>(std::move(self));
  auto me = *this; // 复制任务和参数
  concurrency::create_task([me, self_ptr]() mutable {
    try {
      std::apply(me.fn_, me.args_);
      // 提交到执行器运行complete
      asio::post(self_ptr->get_executor(), [self_ptr]() {
        self_ptr->complete({});
      });
    } catch (std::exception& e) {
      asio::post(self_ptr->get_executor(), [self_ptr, ex = std::current_exception()]() {
        self_ptr->complete(ex);
      });
    }
  });
}

注意:self_ptr要确保在completion handler执行前保持存活,且complete()必须在执行器线程调用,而非task的工作线程。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 06:10:53