【发布时间】:2014-05-24 16:43:29
【问题描述】:
我想探索使用IObservable<T> 作为SqlDataReader 的包装器的可能性。到目前为止,我们一直在使用读取器来避免将整个结果具体化到内存中,并且我们使用阻塞同步 API 来做到这一点。
现在我们想尝试将异步 API 与 .NET Reactive Extensions 结合使用。
但是,由于采用异步方式是一个渐进的过程,因此该代码必须与同步代码共存。
我们已经知道这种同步和异步的混合在 ASP.NET 中不起作用——因为整个请求执行路径必须始终是异步的。一篇关于这个主题的优秀文章是http://blog.stephencleary.com/2012/07/dont-block-on-async-code.html
但我说的是普通的 WCF 服务。我们已经在那里混合了异步和同步代码,但是这是我们第一次想引入 Rx 并且有麻烦。
我创建了简单的单元测试(我们使用 mstest, sigh:-() 来演示问题。我希望有人能够解释我发生了什么。请在下面找到整个源代码(使用起订量):
using System;
using System.Data.Common;
using System.Diagnostics;
using System.Linq;
using System.Reactive.Linq;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.VisualStudio.TestTools.UnitTesting;
using Moq;
namespace UnitTests
{
public static class Extensions
{
public static Task<List<T>> ToListAsync<T>(this IObservable<T> observable)
{
var res = new List<T>();
var tcs = new TaskCompletionSource<List<T>>();
observable.Subscribe(res.Add, e => tcs.TrySetException(e), () => tcs.TrySetResult(res));
return tcs.Task;
}
}
[TestClass]
public class TestRx
{
public const int UNIT_TEST_TIMEOUT = 5000;
private static DbDataReader CreateDataReader(int count = 100, int msWait = 10)
{
var curItemIndex = -1;
var mockDataReader = new Mock<DbDataReader>();
mockDataReader.Setup(o => o.ReadAsync(It.IsAny<CancellationToken>())).Returns<CancellationToken>(ct => Task.Factory.StartNew(() =>
{
Thread.Sleep(msWait);
if (curItemIndex + 1 < count && !ct.IsCancellationRequested)
{
++curItemIndex;
return true;
}
Trace.WriteLine(curItemIndex);
return false;
}));
mockDataReader.Setup(o => o[0]).Returns<int>(_ => curItemIndex);
mockDataReader.CallBase = true;
mockDataReader.Setup(o => o.Close()).Verifiable();
return mockDataReader.Object;
}
private static IObservable<int> GetObservable(DbDataReader reader)
{
return Observable.Create<int>(async (obs, cancellationToken) =>
{
using (reader)
{
while (!cancellationToken.IsCancellationRequested && await reader.ReadAsync(cancellationToken))
{
obs.OnNext((int)reader[0]);
}
}
});
}
[TestMethod, TestCategory("CI"), Timeout(UNIT_TEST_TIMEOUT)]
public void ToListAsyncResult()
{
var reader = CreateDataReader();
var numbers = GetObservable(reader).ToListAsync().Result;
CollectionAssert.AreEqual(Enumerable.Range(0, 100).ToList(), numbers);
Mock.Get(reader).Verify(o => o.Close());
}
[TestMethod, TestCategory("CI"), Timeout(UNIT_TEST_TIMEOUT)]
public void ToEnumerableToList()
{
var reader = CreateDataReader();
var numbers = GetObservable(reader).ToEnumerable().ToList();
CollectionAssert.AreEqual(Enumerable.Range(0, 100).ToList(), numbers);
Mock.Get(reader).Verify(o => o.Close());
}
[TestMethod, TestCategory("CI"), Timeout(UNIT_TEST_TIMEOUT)]
public void ToEnumerableForEach()
{
var reader = CreateDataReader();
int i = 0;
foreach (var n in GetObservable(reader).ToEnumerable())
{
Assert.AreEqual(i, n);
++i;
}
Assert.AreEqual(100, i);
Mock.Get(reader).Verify(o => o.Close());
}
[TestMethod, TestCategory("CI"), Timeout(UNIT_TEST_TIMEOUT)]
public void ToEnumerableForEachBreak()
{
var reader = CreateDataReader();
int i = 0;
foreach (var n in GetObservable(reader).ToEnumerable())
{
Assert.AreEqual(i, n);
++i;
if (i == 5)
{
break;
}
}
Mock.Get(reader).Verify(o => o.Close());
}
[TestMethod, TestCategory("CI"), Timeout(UNIT_TEST_TIMEOUT)]
public void ToEnumerableForEachThrow()
{
var reader = CreateDataReader();
int i = 0;
try
{
foreach (var n in GetObservable(reader).ToEnumerable())
{
Assert.AreEqual(i, n);
++i;
if (i == 5)
{
throw new Exception("xo-xo");
}
}
Assert.Fail();
}
catch (Exception exc)
{
Assert.AreEqual("xo-xo", exc.Message);
Mock.Get(reader).Verify(o => o.Close());
}
}
[TestMethod, TestCategory("CI"), Timeout(UNIT_TEST_TIMEOUT)]
public void Subscribe()
{
var reader = CreateDataReader();
var tcs = new TaskCompletionSource<object>();
int i = 0;
GetObservable(reader).Subscribe(n =>
{
Assert.AreEqual(i, n);
++i;
}, () =>
{
Assert.AreEqual(100, i);
Mock.Get(reader).Verify(o => o.Close());
tcs.TrySetResult(null);
});
tcs.Task.Wait();
}
[TestMethod, TestCategory("CI"), Timeout(UNIT_TEST_TIMEOUT)]
public void SubscribeCancel()
{
var reader = CreateDataReader();
var tcs = new TaskCompletionSource<object>();
var cts = new CancellationTokenSource();
int i = 0;
GetObservable(reader).Subscribe(n =>
{
Assert.AreEqual(i, n);
++i;
if (i == 5)
{
cts.Cancel();
}
}, e =>
{
Assert.IsTrue(i < 100);
Mock.Get(reader).Verify(o => o.Close());
tcs.TrySetException(e);
}, () =>
{
Assert.IsTrue(i < 100);
Mock.Get(reader).Verify(o => o.Close());
tcs.TrySetResult(null);
}, cts.Token);
tcs.Task.Wait();
}
[TestMethod, TestCategory("CI"), Timeout(UNIT_TEST_TIMEOUT)]
public void SubscribeThrow()
{
var reader = CreateDataReader();
var tcs = new TaskCompletionSource<object>();
int i = 0;
GetObservable(reader).Subscribe(n =>
{
Assert.AreEqual(i, n);
++i;
if (i == 5)
{
throw new Exception("xo-xo");
}
}, e =>
{
Assert.AreEqual("xo-xo", e.Message);
Mock.Get(reader).Verify(o => o.Close());
tcs.TrySetResult(null);
});
tcs.Task.Wait();
}
}
}
这些单元测试捕获 API 的所有可能用途,该 API 返回包装数据读取器的 IObservable<T>:
- 人们可能希望使用我们的
ToListAsync扩展方法或.ToEnumerable().ToList()完全实现它。 - 人们可能希望使用
ToEnumerable扩展方法对其进行迭代。是的 - 如果消费很快,它会阻塞,如果消费很慢,它会在内部队列中具体化数据,但这种情况是合法的。 - 最后,人们可以通过订阅来直接使用 observable,但在某些时候他们必须等待结束(阻塞线程),因为周围的大部分代码仍然是同步的。
一个基本要求是数据读取器在读取结束后立即被丢弃 - 无论 observable 以何种方式被消耗。
在所有单元测试中,有 4 个失败:
-
SubscribeCancel和SubscribeThrow超时(即死锁) -
ToEnumerableForEachBreak和ToEnumerableForEachThrow未能验证数据读取器处置。
数据读取器处理验证失败是一个时间问题 - 当foreach 离开(通过异常或中断)时,相应的IEnumerator 立即被处理,最终取消了可观察的实现使用的取消令牌.但是,该实现在另一个线程上运行,并且当它注意到取消时 - 单元测试已经结束。在实际应用程序中,读者会被适当且相当迅速地处理掉,但它不会与迭代结束同步。我想知道是否可以处理上述IEnumerator 实例,直到相应的IObservable 实现注意到取消并且读者被处理掉。
编辑
所以DbDataReader 是IEnumerable,这意味着如果希望同步枚举对象 - 没问题。
但是,如果我想异步执行呢?在这种情况下,我被禁止列举读者——这是一个阻塞操作。唯一的出路是返回一个可观察的。其他人已经用比我更好的语言讨论了这个主题,例如 - http://www.interact-sw.co.uk/iangblog/2013/11/29/async-yield-return
因此我必须返回一个IObservable,而我不能使用ToObservable 扩展方法,因为它取决于阅读器的阻塞枚举。
接下来,给定IObservable,有人可能会将其转换为IEnumerable,这很愚蠢,因为读者已经是IEnumerable,但仍然可行且合法。
编辑 2
使用 .NET Reflector(与 VS 集成)调试代码显示流程通过以下方法:
namespace System.Reactive.Threading.Tasks
{
public static class TaskObservableExtensions
{
...
private static void ToObservableDone<TResult>(Task<TResult> task, AsyncSubject<TResult> subject)
{
switch (task.get_Status())
{
case TaskStatus.RanToCompletion:
subject.OnNext(task.get_Result());
subject.OnCompleted();
return;
case TaskStatus.Canceled:
subject.OnError((Exception) new TaskCanceledException((Task) task));
return;
case TaskStatus.Faulted:
subject.OnError(task.get_Exception().get_InnerException());
return;
}
}
}
}
在异步订阅中取消令牌和从OnNext 抛出都进入此方法(以及成功完成)。取消和抛出都收敛到subject.OnError 方法。该方法应该最终委托给OnError 处理程序。但事实并非如此。
编辑 3
关注Why is the OnError callback never called when throwing from the given subscriber?,我现在想知道满足以下目标的正确方法是什么:
- 通过异步读取
SqlDataReader实例公开可用的对象 - 避免物化对象。实现的选择权应掌握在 API 的调用者手中。
- API 应该可以在异步代码与同步代码混合的环境中使用。为什么?因为我们已经有一台使用同步 IO 的服务器,我们希望逐步淘汰同步阻塞 IO 和异步 IO。
在我面前有这些目标,我想出了这样的东西(请参阅单元测试代码):
private static IObservable<int> GetObservable(DbDataReader reader)
{
return Observable.Create<int>(async (obs, cancellationToken) =>
{
using (reader)
{
while (!cancellationToken.IsCancellationRequested && await reader.ReadAsync(cancellationToken))
{
obs.OnNext((int)reader[0]);
}
}
});
}
这对你有意义吗?如果没有,有什么替代方案?
接下来,我想按照Subscribe 单元测试代码的演示来使用它。但是,SubcribeCancel 和SubscribeThrow 的结果表明这种使用模式是错误的。 Why is the OnError callback never called when throwing from the given subscriber? 解释了错误的原因。
那么,正确的方法是什么?如何防止 API 的使用者错误地使用它(SubcribeCancel 和 SubscribeThrow 是这种错误使用的例子)。
【问题讨论】:
-
附带说明,您的
ToListAsync函数是多余的,因为 Rx 包含以下操作:source.ToList().ToTask();。 -
你可能过度设计你的解决方案...blogs.msdn.com/b/rickandy/archive/2009/11/14/…
-
@ChristopherHarris 有趣的阅读。但是,我们有理由相信我们将从特定应用程序中的异步 IO 中受益。毕竟,我们确实有多个并行的数据库调用。
-
似乎交互式扩展具有 AsyncEnumerable。你可以看看那个。 rx.codeplex.com/SourceControl/latest#Ix.NET/Source/…
-
我会,但它没有解释直接订阅时没有ToEnumerable时的行为。
标签: c# multithreading asynchronous system.reactive sqldatareader