【问题标题】:Using one observable as a clock to test the other for timeout使用一个 observable 作为时钟来测试另一个是否超时
【发布时间】:2019-01-26 04:06:25
【问题描述】:

我想要做的声明如下所示:

// Checks input source for timeouts, based on the number of elements received 
// from clock since the last one received from source. 
// The two selectors are used to generate output elements.
public static IObservable<R> TimeoutDetector<T1,T2,R>(
        this IObservable<T1> source, 
        IObservable<T2> clock, 
        int countForTimeout,
        Func<R> timedOutSelector, 
        Func<T1, R> okSelector)

大理石图在 ascii 中很困难,但这里是:

source --o---o--o-o----o-------------------o---
clock  ----x---x---x---x---x---x---x---x---x---
output --^---^--^-^----^-----------!-------^---

我尝试寻找现有的 Observable 函数,它们可以以我可以使用的方式组合 source 和 clock,但大多数组合函数依赖于接收“每个函数之一”(And , Zip),或者他们从“缺失”的值 (CombineLatest) 中重新返回“上一个”值,或者它们离我需要的值太远了 (Amb, GroupJoin, @987654331 @、Merge、SelectMany、Timeout)。 Sample 看起来很接近,但我不想将源吞吐量限制为时钟速率。

所以现在我被困在试图填补这里的巨大空白:

return new AnonymousObservable<R>(observer =>
{
    //One observer, two observables??
});

抱歉,“您尝试过什么”部分在这里有点弱:假设我已经尝试过思考它!我不是要求完整的实施,只是:

  • 是否有内置功能可以帮助我,但我错过了?
  • 如何构建一个基于 lambda 的观察者,它订阅两个可观察对象?

【问题讨论】:

  • 给未来的读者;在您的问题中,我假设您的大理石图是为 countForTimeout 参数传递 3 的结果?
  • @Lee,是的,没错。

标签: c# system.reactive


【解决方案1】:

我知道您并没有要求完全实施,但我认为这是一个解决方案:

public static IObservable<TR> TimeoutDetector<T1, T2, TR>(
    this IObservable<T1> source,
    IObservable<T2> clock,
    int countForTimeout,
    Func<TR> timedOutSelector,
    Func<T1, TR> okSelector)
{
    return source
        .Select(i => clock.Take(countForTimeout).LastAsync())
        .Switch().Select(_ => timedOutSelector())
        .Merge(source.Select(okSelector));
}

它的工作原理如下 - 我注意到您的输出是 okSelector 投影的源,与超时事件合并。所以我专注于产生超时事件,因为其余的很容易。

这个想法是在每次源发射时创建一个倒计时,并在每个时钟脉冲上递减这个倒计时。如果源发出,我们会中止倒计时,否则当倒计时达到 0 时,我们会产生一个 timedOut 事件。

分解:

  1. 将每个源元素投影到一个采用第一个countForTimeout 元素的流中——注意时钟流必须是“热”可观察的,因为我们在每个 countDown 事件上都订阅它。时钟流变热是很正常的。如果这有事件发生,我们就会超时。
  2. Switch 将丢弃除最新倒计时流之外的所有内容。
  3. 使用Select 投影到timedOut 事件。
  4. 现在只需合并源事件。

这是我使用的单元测试,旨在与您的弹珠图非常相似(nuget rx-testing & nunit 用于编译必要的库):

    [Test]
    public void AKindOfTimeoutTest()
    {
        var scheduler = new TestScheduler();

        var clockStream = scheduler.CreateHotObservable(
            OnNext(100, Unit.Default),
            OnNext(200, Unit.Default),
            OnNext(300, Unit.Default),
            OnNext(400, Unit.Default),
            OnNext(500, Unit.Default),
            OnNext(600, Unit.Default),
            OnNext(750, Unit.Default), /* make clock funky! */
            OnNext(800, Unit.Default),
            OnNext(900, Unit.Default));


        var sourceStream = scheduler.CreateColdObservable(
            OnNext(50, 1),
            OnNext(150, 2),
            OnNext(250, 3),
            OnNext(275, 4),
            OnNext(400, 5),
            OnNext(900, 6));


        Func<int> timedOutSelector = () => 0;
        Func<int, int> okSelector = i => i;

        var results = scheduler.CreateObserver<int>();

        sourceStream.TimeoutDetector(clockStream, 3, timedOutSelector, okSelector)
                    .Subscribe(results);

        scheduler.Start();

        results.Messages.AssertEqual(
            OnNext(50, 1),
            OnNext(150, 2),
            OnNext(250, 3),
            OnNext(275, 4),
            OnNext(400, 5),
            OnNext(750, 0),
            OnNext(900, 6));
    }
}

尝试回答您的具体问题:

  • 问。是否有一个内置功能可以帮助我,我错过了? A. 可能扫描是关键。
  • 问。如何构建一个订阅两个可观察对象的基于 lambda 的观察者? A. 不太清楚你的意思...有很多组合流的方法,你提到了其中的大部分。

【讨论】:

  • 我恨你 :) 我会想出一些更温和的东西......冗长(这也证明了我订阅两个可观察对象的意思)。现在我必须弄清楚你是怎么做的......
  • 我希望你仍然会投票给我,即使你讨厌我。 :)
  • 如果我的source 也很热门,我是否会在您的函数中订阅两次它来冒问题?
  • 如果天气很热,当然不会,有时如果天气很冷,并且您在不同时间订阅但希望事件同步,则有时会出现问题 - 但我在这里不这样做。在寒冷的情况下,您始终可以使用 Publish().RefCount() 组合使其变热/变热。
  • 不错。我要补充一点,您可以通过使用 source.Select(_ =&gt; clock.Take(countForTimeout).LastAsync()) 而不是 startCountdown.Select(i =&gt; clock.Scan(...).Where(...).Take(1)) 来减少一些不必要的流失
【解决方案2】:

这是我提到的 Observable.Create 方法(相同的测试工作):

public static IObservable<TR> TimeoutDetector<T1, T2, TR>(
    this IObservable<T1> source,
    IObservable<T2> clock,
    int countForTimeout,
    Func<TR> timedOutSelector,
    Func<T1, TR> okSelector)
{
    return Observable.Create<TR>(observer =>
        {
            var counter = countForTimeout;

            var timeoutSub = clock.Subscribe(_ =>
                {
                    var count = Interlocked.Decrement(ref counter);
                    if (count == 0)
                    {
                        observer.OnNext(timedOutSelector());
                    }
                },
                observer.OnError,
                observer.OnCompleted);

            var sourceSub = source.Subscribe(
                i =>
                {
                    Interlocked.Exchange(ref counter, countForTimeout);
                    observer.OnNext(okSelector(i));
                },
                observer.OnError,
                observer.OnCompleted);

            return new CompositeDisposable(sourceSub, timeoutSub);
        });
}

请注意,Observable.Create 对确保使用正确的 Rx 语法非常有帮助(即流发出 OnNext* (OnError | OnCompleted)? - 这意味着我可以稍微放松一下最多发送一次 OnError 或 OnCompleted。

【讨论】:

  • 很好,我想知道 OnError 和 OnCompleted 中所有那些看似多余的锁和检查。
  • 我可以再次投票。我刚刚切换回这个实现,因为我意识到接受的答案没有检测到源 never 是否触发(并且它也会重复订阅/取消订阅时钟,虽然我是不知道为什么)。
【解决方案3】:

我想出了这个,这比詹姆斯的回答要漂亮。

public static IObservable<R> TimeoutDetector2<T1, T2, R>(
        this IObservable<T1> source, 
        IObservable<T2> clock, int maxDiff,
        Func<R> timedOutSelector, Func<T1, R> okSelector)
{
    return new AnonymousObservable<R>(observer =>
    {
        int counter = 0;
        object gate = new object();
        bool error = false;
        bool completed = false;
        bool timedOut = false;
        var sourceSubscription = source.Subscribe(
            x =>
            {
                lock(gate)
                {
                    if(!error && !completed) observer.OnNext(okSelector(x));
                    counter = 0;
                    timedOut = false;
                }
            },
            ex =>
            {
                lock(gate)
                {
                    error = true;
                    if(!completed) observer.OnError(ex);
                }
            },
            () =>
            {
                lock(gate)
                {
                    completed = true;
                    if(!error) observer.OnCompleted();
                }
            });
        var clockSubscription = clock.Subscribe(
            x =>
            {
                lock(gate)
                {
                    counter = counter + 1;
                    if(!error && !completed && counter > maxDiff && !timedOut)
                    {
                        timedOut = true;
                        observer.OnNext(timedOutSelector());
                    }
                }
            },
            ex =>
            {
                lock(gate)
                {
                    error = true;
                    if(!completed) observer.OnError(ex);
                }
            },
            () =>
            {
                lock(gate)
                {
                    completed = true;
                    if(!error) observer.OnCompleted();
                }
            });

        //need to return a subscription
        return new CompositeDisposable(sourceSubscription, clockSubscription);
    }).Publish().RefCount(); // prevent subscribers provoking more than one subscription to source and clock
}

【讨论】:

  • 旋转它并使用 Observable.Create 而不是创建 AnonymousObservable - 我认为这对你来说会容易得多。您可以在 Create 中订阅您喜欢的任何内容,并从 System.Reactive.Disposables 命名空间返回一个时髦的 IDisposable(如 CompositeDisposable)以进行清理。
  • 嗯,也许我没有看到它,但我假设只是用 Observable.Create&lt;R&gt; 替换新的 AnonymousObservable&lt;R&gt; 不是你的意思......
  • 否,使用下面的Create 方法添加了答案。希望您能看到它允许使用与您所采用的相同类型的命令式方法,但它构建起来更干净,多次订阅也更安全。
  • 它不再是“下面”:)
  • 我总是这样做。呸!在“附近某处”查看其他答案! :)
【解决方案4】:

当然,这是一个老问题。我一直在寻找比timeoutWith 更高级的东西,以便我可以取消超时逻辑。

我想知道超时是否几乎等同于:

race(throwError('timedout').pipe(delay(10000)), yourObs$)

那么这里显示的“throwError”当然是可以取消的。

如果您想知道为什么 - 我有一些由可观察链控制的“步骤”并且我有一个超时。但是,如果其中一个步骤包括打开一个对话框,那么我希望取消超时!

【讨论】:

  • 我猜timeout 在你写这个问题时可能不存在。也许race 也没有!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-10-21
  • 1970-01-01
  • 1970-01-01
  • 2020-11-17
  • 1970-01-01
  • 1970-01-01
  • 2019-01-23
相关资源
最近更新 更多