【问题标题】:How to wrap SqlDataReader with IObservable properly?如何用 IObservable 正确包装 SqlDataReader?
【发布时间】: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&lt;T&gt;

  • 人们可能希望使用我们的ToListAsync 扩展方法或.ToEnumerable().ToList() 完全实现它。
  • 人们可能希望使用ToEnumerable 扩展方法对其进行迭代。是的 - 如果消费很快,它会阻塞,如果消费很慢,它会在内部队列中具体化数据,但这种情况是合法的。
  • 最后,人们可以通过订阅来直接使用 observable,但在某些时候他们必须等待结束(阻塞线程),因为周围的大部分代码仍然是同步的。

一个基本要求是数据读取器在读取结束后立即被丢弃 - 无论 observable 以何种方式被消耗。

在所有单元测试中,有 4 个失败:

  • SubscribeCancelSubscribeThrow 超时(即死锁)
  • ToEnumerableForEachBreakToEnumerableForEachThrow 未能验证数据读取器处置。

数据读取器处理验证失败是一个时间问题 - 当foreach 离开(通过异常或中断)时,相应的IEnumerator 立即被处理,最终取消了可观察的实现使用的取消令牌.但是,该实现在另一个线程上运行,并且当它注意到取消时 - 单元测试已经结束。在实际应用程序中,读者会被适当且相当迅速地处理掉,但它不会与迭代结束同步。我想知道是否可以处理上述IEnumerator 实例,直到相应的IObservable 实现注意到取消并且读者被处理掉。

编辑

所以DbDataReaderIEnumerable,这意味着如果希望同步枚举对象 - 没问题。

但是,如果我想异步执行呢?在这种情况下,我被禁止列举读者——这是一个阻塞操作。唯一的出路是返回一个可观察的。其他人已经用比我更好的语言讨论了这个主题,例如 - 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?,我现在想知道满足以下目标的正确方法是什么:

  1. 通过异步读取SqlDataReader 实例公开可用的对象
  2. 避免物化对象。实现的选择权应掌握在 API 的调用者手中。
  3. 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 单元测试代码的演示来使用它。但是,SubcribeCancelSubscribeThrow 的结果表明这种使用模式是错误的。 Why is the OnError callback never called when throwing from the given subscriber? 解释了错误的原因。

那么,正确的方法是什么?如何防止 API 的使用者错误地使用它(SubcribeCancelSubscribeThrow 是这种错误使用的例子)。

【问题讨论】:

  • 附带说明,您的 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


【解决方案1】:

订阅取消

SubscribeCancelcts 取消而失败。这不会调用 OnError 处理程序。

取消您的cts 等同于处置您的订阅。处理订阅会导致所有未来的OnNextOnErrorOnCompleted 调用被忽略。因此,任务永远不会完成,测试永远挂起。

解决方案:

当您取消 cts 时,将任务设置为正确的状态。

订阅投掷

SubscribeThrow 由于OnNext 处理程序中的异常而失败。

OnNext 处理程序中引发异常不会将异常转发到OnError 处理程序。

解决方案:

不要在您的 Subscribe 处理程序中抛出异常。相反,请处理您的订阅并将Task 设置为正确的状态。

ToEnumerableForEachThrow & ToEnumerableForEachBreak

ToEnumerableForEachThrowToEnumerableForEachBreak 由于竞争条件而失败。

foreach(...) 上的 enumerable 将调用底层 observable 上的 dispose,这将取消取消令牌。之后,异常被您的测试的 catch 捕获(或者中断只是退出 foreach),您可以在其中测试以查看底层阅读器已被释放......除非阅读器尚未被释放,因为 observable 仍然等待读者产生下一个结果......只有在读者产生(和可观察的产生)之后,可观察的循环才会返回并检查取消令牌。此时,可观察对象中断并退出 using 块并处置阅读器。

解决方案:

从您的Observable.Create 返回一个Disposable,而不是您的using (...) 语句。一次性将在订阅被处置时被处置。这就是你想要的。把using 语句全部去掉,让Rx 完成它的工作。

【讨论】:

  • 我不同意SubscribeCancel 的说法。令牌可能在代码的深处被取消。将相应的 Task 实例传播到可能取消令牌的所有地方的要求是疯狂的,并且不支持使用 .NET Reflector 检查 Rx 代码 - 我将立即在我的帖子中添加一个编辑。
  • 您不同意哪一部分?事实陈述,还是我推荐的解决方案?
  • 请参阅我在帖子中的编辑,其中我声称 Rx 代码以完全相同的方式处理令牌的取消和从 OnNext 抛出,并且预计两者都会导致 OnError ,但出于某种原因不要这样做。
  • 向我显示将其定义为设计行为的文档的链接。
  • 让我们通过聊天继续...chat.stackoverflow.com/rooms/54553/rx-behavior
猜你喜欢
  • 2013-08-10
  • 2017-06-15
  • 1970-01-01
  • 2014-05-07
  • 2016-12-23
  • 1970-01-01
  • 2016-04-16
  • 2017-11-30
  • 1970-01-01
相关资源
最近更新 更多