【问题标题】:Ensuring completion of async OnNext code before process terminates确保在进程终止之前完成异步 OnNext 代码
【发布时间】:2016-09-13 22:10:45
【问题描述】:

下面的单元测试永远不会打印“Async 3”,因为测试首先完成。我怎样才能确保它运行完成?我能想到的最好的方法是最后的任意 Task.Delay 或 WriteAsync().Result,两者都不理想。

    public async Task TestMethod1() // eg. webjob
    {
        TestContext.WriteLine("Starting test...");
        var observable = Observable.Create<int>(async ob =>
        {
            ob.OnNext(1);
            await Task.Delay(1000); // Fake async REST api call
            ob.OnNext(2);
            await Task.Delay(1000);
            ob.OnNext(3);
            ob.OnCompleted();
        });

        observable.Subscribe(i => TestContext.WriteLine($"Sync {i}"));
        observable.SelectMany(i => WriteAsync(i).ToObservable()).Subscribe();

        await observable;
        TestContext.WriteLine("Complete.");
    }

    public async Task WriteAsync(int value) // Fake async DB call
    {
        await Task.Delay(1000);
        TestContext.WriteLine($"Async {value}");
    }

编辑 我意识到提到单元测试可能会产生误导。这不是一个测试问题。该代码模拟了 Azure WebJob 中运行的进程的实际问题,其中生产者和消费者都需要调用一些 Async IO。问题是网络作业在消费者真正完成之前运行完成。这是因为我无法弄清楚如何正确地等待消费者方面的任何事情。也许这对 RX 来说是不可能的......

【问题讨论】:

  • 你知道你正在为底层的 observable 创建三个独立的订阅吗?这是你的意图吗?或者您是否尝试在三个观察者之间共享单个 observable 的值?
  • 3 个订阅? 'await observable' 是否也会导致订阅?在实际代码中,原始 observable 很热,所以我可能应该在示例中也让它变得很热......
  • 是的,await observable 确实会导致第三次订阅。您的第二个订阅 observable.SelectMany(i =&gt; WriteAsync(i).ToObservable()).Subscribe() 有延迟,因此 await observable 在它之前完成。这就是为什么你错过了"Async 3"

标签: c# system.reactive


【解决方案1】:

编辑: 您基本上是在寻找阻塞运算符。旧的阻塞操作符(如ForEach)已被弃用,取而代之的是异步版本。您想像这样等待最后一项:

public async Task TestMethod1()
{
    TestContext.WriteLine("Starting test...");
    var observable = Observable.Create<int>(async ob =>
    {
        ob.OnNext(1);
        await Task.Delay(1000);
        ob.OnNext(2);
        await Task.Delay(1000);
        ob.OnNext(3);
        ob.OnCompleted();
    });

    observable.Subscribe(i => TestContext.WriteLine($"Sync {i}"));
    var selectManyObservable = observable.SelectMany(i => WriteAsync(i).ToObservable()).Publish().RefCount();
    selectManyObservable.Subscribe();
    await selectManyObservable.LastOrDefaultAsync();
    TestContext.WriteLine("Complete.");
}

虽然这将解决您当前的问题,但由于以下原因,您似乎会继续遇到问题(我又添加了两个)。正确使用时 Rx 非常强大,使用不当时会令人困惑。

旧答案:


几件事:

  1. 混合使用 async/await 和 Rx 通常会导致两者兼而有之,而两者都没有好处。
  2. Rx 具有强大的测试功能。你没有使用它。
  3. 副作用,如 WriteLine,最好只在订阅中执行,而不是在像 SelectMany 这样的运算符中执行。
  4. 您可能想复习一下冷与热 observables。
  5. 它没有运行完成的原因是因为您的测试运行器。您的测试运行程序将在TestMethod1 结束时终止测试。否则,Rx 订阅将继续存在。当我在 Linqpad 中运行您的代码时,我得到以下输出:

    开始测试...
    同步 1
    同步 2
    异步 1
    同步 3
    异步 2
    完成。
    异步 3

...这是我假设您想看到的,除了您可能想要 Async 3 之后的 Complete。


仅使用 Rx,您的代码将如下所示:

public void TestMethod1()
{
    TestContext.WriteLine("Starting test...");
    var observable = Observable.Concat<int>(
        Observable.Return(1),
        Observable.Empty<int>().Delay(TimeSpan.FromSeconds(1)),
        Observable.Return(2),
        Observable.Empty<int>().Delay(TimeSpan.FromSeconds(1)),
        Observable.Return(3)
    );

    var syncOutput = observable
        .Select(i => $"Sync {i}");
    syncOutput.Subscribe(s => TestContext.WriteLine(s));

    var asyncOutput = observable
        .SelectMany(i => WriteAsync(i, scheduler));
    asyncOutput.Subscribe(s => TestContext.WriteLine(s), () => TestContext.WriteLine("Complete."));
}

public IObservable<string> WriteAsync(int value, IScheduler scheduler)
{
    return Observable.Return(value)
        .Delay(TimeSpan.FromSeconds(1), scheduler)
        .Select(i => $"Async {value}");
}


public static class TestContext
{
    public static void WriteLine(string s)
    {
        Console.WriteLine(s);
    }
}

这仍然没有利用 Rx 的测试功能。看起来像这样:

public void TestMethod1()
{
    var scheduler = new TestScheduler();
    TestContext.WriteLine("Starting test...");
    var observable = Observable.Concat<int>(
        Observable.Return(1),
        Observable.Empty<int>().Delay(TimeSpan.FromSeconds(1), scheduler),
        Observable.Return(2),
        Observable.Empty<int>().Delay(TimeSpan.FromSeconds(1), scheduler),
        Observable.Return(3)
    );

    var syncOutput = observable
        .Select(i => $"Sync {i}");
    syncOutput.Subscribe(s => TestContext.WriteLine(s));

    var asyncOutput = observable
        .SelectMany(i => WriteAsync(i, scheduler));
    asyncOutput.Subscribe(s => TestContext.WriteLine(s), () => TestContext.WriteLine("Complete."));

    var asyncExpected = scheduler.CreateColdObservable<string>(
        ReactiveTest.OnNext(1000.Ms(), "Async 1"),
        ReactiveTest.OnNext(2000.Ms(), "Async 2"),
        ReactiveTest.OnNext(3000.Ms(), "Async 3"),
        ReactiveTest.OnCompleted<string>(3000.Ms() + 1) //+1 because you can't have two notifications on same tick
    );

    var syncExpected = scheduler.CreateColdObservable<string>(
        ReactiveTest.OnNext(0000.Ms(), "Sync 1"),
        ReactiveTest.OnNext(1000.Ms(), "Sync 2"),
        ReactiveTest.OnNext(2000.Ms(), "Sync 3"),
        ReactiveTest.OnCompleted<string>(2000.Ms()) //why no +1 here?
    );

    var asyncObserver = scheduler.CreateObserver<string>();
    asyncOutput.Subscribe(asyncObserver);
    var syncObserver = scheduler.CreateObserver<string>();
    syncOutput.Subscribe(syncObserver);
    scheduler.Start();
    ReactiveAssert.AreElementsEqual(
        asyncExpected.Messages,
        asyncObserver.Messages);

    ReactiveAssert.AreElementsEqual(
        syncExpected.Messages,
        syncObserver.Messages);
}

public static class MyExtensions
{
    public static long Ms(this int ms)
    {
        return TimeSpan.FromMilliseconds(ms).Ticks;
    }

}

...因此,与您的任务测试不同,您不必等待。测试立即执行。您可以将延迟时间提高到几分钟或几小时,TestScheduler 基本上会为您模拟时间。然后你的测试运行者可能会很高兴。

【讨论】:

  • 如果有人可以在 +1 上填写我的信息,我将非常感激。
  • 我在问题中添加了更多细节,这可能有点模棱两可——这并不是关于 RX 测试的问题,尽管这是我打算进一步了解的内容
  • 谢谢。编辑答案。
【解决方案2】:

好吧,您可以使用Observable.ForEach 阻止直到IObservable 终止:

observable.ForEach(unusedValue => { });

你能不能把TestMethod1变成一个普通的非异步方法,然后用这个替换await observable;

【讨论】:

  • 不,不幸的是它产生了相同的输出(ForEach 也被弃用了)。
  • OP 已经在await observable 上阻塞,但第二个订阅独立运行,因为它被延迟,所以它稍后完成。
猜你喜欢
  • 2017-03-17
  • 1970-01-01
  • 1970-01-01
  • 2015-12-12
  • 1970-01-01
  • 1970-01-01
  • 2016-01-12
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多