【问题标题】:System.Threading.Channels ReadAsync() method is blocking executionSystem.Threading.Channels ReadAsync() 方法正在阻塞执行
【发布时间】:2020-08-03 17:17:06
【问题描述】:

概述

我正在尝试围绕IObserver<T> 接口编写IAsyncEnumerable<T> 包装器。起初我使用BufferBlock<T> 作为后备数据存储,但通过性能测试和研究发现它实际上是一个非常慢的类型,所以我决定试一试System.Threading.Channels.Channel 类型。我的 BufferBlock 实现遇到了与这个类似的问题,但这次我不确定如何解决。

问题

如果我的IObserver<T>.OnNext() 方法尚未写入_channel,我的GetAsyncEnumerator() 循环会被await _channel.Reader.WaitToRead(token) 调用阻塞。在不阻塞程序执行的情况下,在这种情况下等待一个可用的值的正确方法是什么?

实施

public sealed class ObserverAsyncEnumerableWrapper<T> : IAsyncEnumerable<T>, IObserver<T>, IDisposable
{
    private readonly IDisposable _unsubscriber;
    private readonly Channel<T> _channel = Channel.CreateUnbounded<T>();

    private bool _producerComplete;

    public ObserverAsyncEnumerableWrapper(IObservable<T> provider)
    {
        _unsubscriber = provider.Subscribe(this);
    }

    public async void OnNext(T value)
    {
        Log.Logger.Verbose("Adding value to Channel.");
        await _channel.Writer.WriteAsync(value);
    }

    public void OnError(Exception error)
    {
        _channel.Writer.Complete(error);
    }

    public void OnCompleted()
    {
        _producerComplete = true;
    }

    public async IAsyncEnumerator<T> GetAsyncEnumerator([EnumeratorCancellation] CancellationToken token = new CancellationToken())
    {
        Log.Logger.Verbose("Starting async iteration...");
        while (await _channel.Reader.WaitToReadAsync(token) || !_producerComplete)
        {
            Log.Logger.Verbose("Reading...");
            while (_channel.Reader.TryRead(out var item))
            {
                Log.Logger.Verbose("Yielding item.");
                yield return item;
            }
            Log.Logger.Verbose("Awaiting more items.");
        }
        Log.Logger.Verbose("Iteration Complete.");
        _channel.Writer.Complete();
    }

    public void Dispose()
    {
        _channel.Writer.Complete();
        _unsubscriber?.Dispose();
    }
}

附加上下文

没关系,但在运行时传递给构造函数的IObservable&lt;T&gt; 实例是从对Microsoft.Management.Infrastructure api 进行的异步调用返回的CimAsyncResult。这些使用了我试图用花哨的新异步枚举模式包装的 Observer 设计模式。

编辑

更新了对调试器输出的日志记录,并按照一位评论者的建议使我的 OnNext() 方法异步/等待。你可以看到它永远不会进入while()循环。

【问题讨论】:

  • 微妙的评论:你从来没有await或检查来自_channel.Writer.WriteAsync(value);的回复
  • 我按照您的建议更新了 OnNext() 方法,但没有看到效果。
  • 我应该提一下,我不相信 OnNext() 在任何情况下都会异步执行,因为调用我的 OnNext 方法的 IObervable 也不是异步的,并且是第 3 方库类型。话虽如此,我认为这不会影响我的枚举员的行为。
  • 好的。所以我弄清楚了为什么我的程序会阻止执行。在调用堆栈的更远处,我通过GetAwaiter().GetResult() 方法同步调用异步方法。我发现即使使用 Reactive 的 ToAsyncEnumerable() 扩展方法,我得到的结果也是一样的,即程序被阻止了。我这样做是因为有一次我想从构造函数中获取数据。我更改了该实现以使用 Task.Run() 执行调用,现在迭代器可以在两种实现中完美运行。
  • 我查看了来自 Reactive 的 ToAsyncEnumerable 实现,它的构建远远超过了我的实现。它使用ConcurrentQueue&lt;&gt; 作为其数据源。我不知道性能概况与Channel 相比如何,但我认为这很有趣。

标签: asynchronous .net-core iasyncenumerable


【解决方案1】:

在调用堆栈的上方,我通过 GetAwaiter().GetResult() 方法同步调用异步方法。

是的,that's a problem

我这样做是因为有一次我想从构造函数中获取数据。我更改了该实现以使用 Task.Run() 执行调用,现在迭代器可以在两种实现中完美运行。

有比阻塞异步代码更好的解决方案。 Using Task.Run is one way to avoid the deadlock,但您最终仍会获得低于标准的用户体验(我假设您的用户体验是一个 UI 应用程序,因为有一个 SynchronizationContext)。

如果使用异步枚举器加载数据进行显示,那么更合适的解决方案是(synchronously) initialize the UI to a "Loading..." state, and then update that state as the data is loaded asynchronously。如果异步枚举器用于其他用途,您可能会在我的async constructors blog post 中找到一些合适的替代模式。

【讨论】:

  • “Stephen Cleary 回答了你的问题”有成就吗? :D 感谢您提供博客文章的链接,我一定会阅读的。就像我在我的一个 cmets 中提到的那样,我最初的“GetAsyncEnumerator”实现在 while 循环中没有 await 关键字。相反,它等待配置为一次只允许访问一个方法的“SemaphoreSlim.WaitAsync()”。当我将实现更改为我的问题时,在主线程上执行它会导致它停止并且不允许我的生产者插入任何项目。
猜你喜欢
  • 1970-01-01
  • 2011-01-07
  • 1970-01-01
  • 2010-09-08
  • 1970-01-01
  • 1970-01-01
  • 2012-07-13
  • 1970-01-01
  • 2020-02-26
相关资源
最近更新 更多