【问题标题】:Wait async for a event from a never ending Task等待来自永无止境任务的事件的异步
【发布时间】:2020-01-07 08:48:20
【问题描述】:

目前我正在构建一个需要等待传入消息的网络服务器,但我需要异步等待这些消息。这些消息在另一个继续运行的任务中接收。无休止的任务推送事件,另一个任务也需要等待这些事件之一被接收。

如果需要的话会画一个草图来更好地可视化问题

我已经看到了 TaskCompletionSource 类,但除非客户端断开连接,否则该任务将无法完成,这样就无法工作。等待的任务完全一样。

是否有适用于 C#/.net 核心的库或内置解决方案?

谢谢你:)

【问题讨论】:

  • 到目前为止您尝试过什么?如果这是您真正想要的,您应该能够使用 TaskTaskCompletionSource 创建此行为。 stackoverflow.com/a/29355817/10608418
  • 我不确定我是否了解您的需求。 其他任务需要等待这些事件之一是什么意思?异步等待事件听起来像是一个正常的事件。该事件在永无止境的任务的范围内被调用。你的意思是这样的:DEMO?

标签: c# events .net-core channel system.threading.channels


【解决方案1】:

.NET Core 2.1 和 System.Threading.Channels 中有新的 Channel<T> API,可让您异步使用队列:

Channel<int> channel = Channel.CreateUnbounded<int>();
...
ChannelReader<int> c = channel.Reader;
while (await c.WaitToReadAsync())
{
    if (await c.ReadAsync(out int item))
    {
         // process item...
    }
}

请参考this博客文章了解如何使用它的介绍和更多示例。

【讨论】:

  • Channels 已经在 SignalR 中用于实现event streaming。那里显示的模式很重要 - 在生产者方法内部创建通道,并且只公开 ChannelReader。这样就不会混淆谁可以写入频道或谁负责关闭它
【解决方案2】:

频道

ASP.NET Core SignalR 使用 Channels to implement event streaming - 发布者方法异步生成由另一个方法处理的事件。在 SignalR 的情况下,它是一个将每个新事件推送给客户端的消费者。

在您的代码中执行相同的操作很容易,但请确保遵循 SignalR 文档中显示的模式 - 通道由工作人员自己创建,并且永远不会暴露给调用者。这意味着通道的状态没有歧义,因为只有工作人员可以关闭它。

另一个重要的一点是你必须在完成后关闭通道,即使有异常。否则,消费者将无限期地阻塞。

您还需要将 CancellationToken 传递给生产者,以便您可以在某个时候终止它。即使您不打算显式取消生产者,您也需要一种方法来告诉它在应用程序终止时停止。

以下方法创建通道确保即使发生异常也会关闭。

ChannelReader<int> MyProducer(someparameters,CancellationToken token)
{
    var channel=Channel.CreateUnbounded<int>();
    var writer=channel.Writer;

    _ = Task.Run(async()=>{
            while(!token.IsCancellationRequested)
            {
                var i= ... //do something to produce a value
                await writer.WriteAsync(i,token);
            }
        },token)
        //IMPORTANT: Close the channel no matter what.
        .ContinueWith(t=>writer.Complete(t.Exception));                
    return channel.Reader;
}

t.Exception 如果生产者正常完成或通过令牌取消工作任务,则为 null。没有理由使用ChannelReader.TryComplete,因为一次只有一个任务是写给作者。

消费者只需要那个阅读器来消费事件:

async Task MyConsumer(ChannerReader<int> reader,CancellationToken token)
{
    while (await reader.WaitToReadAsync(cts.Token).ConfigureAwait(false))
    {
        while (reader.TryRead(out int item))
        {
           //Use that integer
        }
    }
}

可用于:

var reader=MyProducer(...,cts.Token);
await MyConsumer(reader,cts.Token);

IAsyncEnumerable

我在这里作弊了,因为我从 ChannelReader.ReadAllAsync 方法中复制了循环代码,该方法返回一个允许异步循环的 IAsyncEnumerable。在 .NET Core 3.0 和 .NET Standard 2.1(但不仅如此)中,可以将使用者替换为:

await foreach (var i from reader.ReadAllAsync(token))
{
     //Use that integer
}

Microsoft.Bcl.AsyncInterfaces 包向以前的 .NET 版本添加了相同的接口

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-03-18
    • 1970-01-01
    • 2014-09-06
    • 1970-01-01
    • 2013-02-10
    • 2017-07-22
    相关资源
    最近更新 更多