【发布时间】: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<T> 实例是从对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<>作为其数据源。我不知道性能概况与Channel相比如何,但我认为这很有趣。
标签: asynchronous .net-core iasyncenumerable