【问题标题】:How to compose asynchronous operations?如何组合异步操作?
【发布时间】:2019-08-25 09:15:52
【问题描述】:

我正在寻找一种组合异步操作的方法。最终目标是执行异步操作,或者让它运行完成,或者在用户定义的超时后返回。

出于示例目的,假设我正在寻找一种方法来组合以下协程1

IAsyncOperation<IBuffer> read(IBuffer buffer, uint32_t count)
{
    auto&& result{ co_await socket_.InputStream().ReadAsync(buffer, count, InputStreamOptions::None) };
    co_return result;
}

socket_StreamSocket 实例。

还有超时协程:

IAsyncAction timeout()
{
    co_await 5s;
}

我正在寻找一种方法来组合这些协程,以便在读取数据或超时后尽快返回。

这些是我目前评估过的选项:

  • C++20 协程:据我了解P1056R0,目前没有库或语言功能“能够创建和组合协程”
  • Windows 运行时提供了异步任务类型,最终派生自 IAsyncInfo:同样,我没有找到任何工具可以让我以我需要的方式组合任务。
  • Concurrency Runtime:这看起来很有希望,尤其是 when_any 函数模板看起来正是我所需要的。

看来我需要使用并发运行时。但是,我很难将所有部分组合在一起。我对如何处理异常以及是否需要取消相应的其他并发任务感到特别困惑。

问题有两个:

  • 并发运行时是唯一的选项(UWP 应用程序)吗?
  • 实现是什么样的?

1方法在应用程序内部。不需要让它们返回与 Windows 运行时兼容的类型。

【问题讨论】:

    标签: c++ uwp windows-runtime c++-winrt concurrency-runtime


    【解决方案1】:

    我认为最简单的方法是使用concurrency 库。您需要修改超时以返回与第一种方法相同的类型,即使它返回 null。

    (我意识到这只是部分答案......)

    我的 C++ 很烂,但我认为这很接近......

    array<task<IBuffer>, 2> tasks =
    {
    concurrency::create_task([]{return read(buffer, count).get();}),
    concurrency::create_task([]{return modifiedTimeout.get();})
    };
    
    concurrency::when_any(begin(tasks), end(tasks)).then([](IBuffer buffer)
    { 
        //do something 
    });
    
    

    【讨论】:

    • 这看起来很有希望,尽管它并没有计算出与取消和异常传播有关的细节。选择std::optional 作为task 的返回类型看起来像是一个有效的选项。请注意,when_any 返回一个std::pair&lt;T, size_t&gt;,因此需要在then 延续中使用。我看看能不能用你提供的信息算出剩下的细节。
    • 既然你在做基于事件的东西,你可能想要看看的替代方案是 ReactiveX。它有一个 C++ 库。但是,您基本上为您的 IBuffer 创建了一个 Observable,并为您的计时器创建了一个 Observable(它也发出 IBuffer,但为空)。然后合并它们并 Take(1)。这将“完成”合并的 Observable,之后您可以取消尚未完成的 Observable。异常也很容易捕获,尽管我从未尝试过使用 C++,只有 C#。如果您从未使用过 ReactiveX 中的 Observables,请注意学习曲线陡峭
    • 终于用你提供的信息做了一些测试。还有一些开放式的结局,我还没有完全理解。我原本打算更新你的答案,但结果比我想象的要长得多,我不想把你的答案变成你可能不同意的东西。 ReactiveX 看起来很有用,尽管它对于我(当前)需要的东西来说有点太多了。而且我对添加另一个基于模板的库并不感到兴奋。几乎不可能像现在这样对 C++/WinRT 项目定期运行代码分析...
    【解决方案2】:

    正如 Lee McPherson 在另一个 answer 中所建议的那样,Concurrency Runtime 看起来是一个可行的选择。它提供了tasks,可以与其他人组合,使用延续链接,以及与 Windows 运行时异步模型无缝集成(请参阅Creating Asynchronous Operations in C++ for UWP Apps)。作为奖励,包括 &lt;pplawait.h&gt; 标头为 concurrency::task 类模板实例提供了适配器,可用作 C++20 协程等待对象。

    我无法回答所有问题,但这是我最终想出的。为简单起见(和便于验证),我使用Sleep 代替实际的读取操作,并返回int 而不是IBuffer

    任务组成

    ConcRT 提供了多种组合任务的方法。鉴于要求concurrency::when_any 可用于创建一个任务,该任务在任何提供的任务完成时返回。当仅提供 2 个任务作为输入时,还有一个便利运算符 (operator||) 可用。

    异常传播

    任一输入任务引发的异常都不算作成功完成。当与when_any 任务一起使用时,抛出异常不足以满足等待条件。因此,异常不能用于中断组合任务。为了解决这个问题,我选择返回 std::optional,并在 then 延续中引发适当的异常。

    任务取消

    这对我来说仍然是个谜。看来,一旦任务满足when_any 任务的等待条件,就不需要取消相应的其他未完成任务。一旦完成(成功或失败),它们就会被默默地处理。

    以下是代码,使用前面提到的简化。它创建了一个由实际工作负载和超时任务组成的任务,两者都返回一个std::optionalthen 延续检查返回值,并在没有返回值的情况下抛出异常(即 timeout_task 首先完成)。

    #include <Windows.h>
    
    #include <cstdint>
    #include <iostream>
    #include <optional>
    #include <ppltasks.h>
    #include <stdexcept>
    
    using namespace concurrency;
    
    task<int> read_with_timeout(uint32_t read_duration, uint32_t timeout)
    {
        auto&& read_task
        {
            create_task([read_duration]
                {
                    ::Sleep(read_duration);
                    return std::optional<int>{42};
                })
        };
        auto&& timeout_task
        {
            create_task([timeout]
                {
                    ::Sleep(timeout);
                    return std::optional<int>{};
                })
        };
    
        auto&& task
        {
            (read_task || timeout_task)
            .then([](std::optional<int> result)
                {
                    if (!result.has_value())
                    {
                        throw std::runtime_error("timeout");
                    }
                    return result.value();
                })
        };
        return task;
    }
    

    以下测试代码

    int main()
    {
        try
        {
            auto res1{ read_with_timeout(3000, 5000).get() };
            std::cout << "Succeeded. Result = " << res1 << std::endl;
            auto res2{ read_with_timeout(5000, 3000).get() };
            std::cout << "Succeeded. Result = " << res2 << std::endl;
        }
        catch( std::runtime_error const& e )
        {
            std::cout << "Failed. Exception = " << e.what() << std::endl;
        }
    }
    

    产生这个输出:

    Succeeded. Result = 42
    Failed. Exception = timeout
    

    【讨论】:

      猜你喜欢
      • 2020-05-01
      • 2023-03-16
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-07-03
      • 1970-01-01
      • 2017-07-02
      相关资源
      最近更新 更多