【问题标题】:Rx.NET Take an element and subscribe again after some timeRx.NET 获取一个元素并在一段时间后再次订阅
【发布时间】:2018-11-14 00:40:35
【问题描述】:

我需要一种节流,它的工作方式有点不同。我需要从序列中获取一个元素,取消订阅并在 1 秒内再次订阅。换句话说,我想在获取第一个元素后的 1 秒内忽略所有元素:

Input:  (1) -100ms- (2) -200ms- (3) -1_500ms- (4) -1_000ms- (5) -500ms- (6) ...
Output: (1) --------------------------------- (4) --------- (5) ----------- ...

如何用 Rx.NET 实现这个简单的事情?

【问题讨论】:

    标签: reactive-programming system.reactive rx.net


    【解决方案1】:

    试试这个:

    Input
        .Window(() => Observable.Timer(TimeSpan.FromSeconds(1.0)))
        .SelectMany(xs => xs.Take(1));
    

    这是一个测试:

    var query =
        Observable
            .Interval(TimeSpan.FromSeconds(0.2))
            .Window(() => Observable.Timer(TimeSpan.FromSeconds(1.0)))
            .SelectMany(xs => xs.Take(1));
    

    这产生了:

    0 5 10 14 19 24 29 34 39

    从 10 跳转到 14 只是使用多线程的结果,而不是查询中的错误。

    【讨论】:

    • @Eugene - 我测试了它,它工作正常。你的来源是什么?
    • @Eugene 我也试过这段代码,我每 1 秒才收到一次事件。
    • 我的情况是:第一个元素可以在开始后 0.8 秒后出现,而不是在开始后(如您的示例中)。而且我需要立即获取第一个元素(+重新启动 1 秒锁定期),而不是在第一秒结束时。
    • @Enigmativity 看起来你的例子工作正常,我做了更多的实验。
    • 这个答案并不完全符合规范,尽管它可能适用于您的情况。查看替代答案。
    【解决方案2】:

    @Enigmativity 的回答并不完全符合规范。它可能适用于你想要的东西。

    他的答案定义了 1 秒窗口,并从每个窗口中获取第一个窗口。但是,这并不能保证您在项目之间保持一秒钟的沉默。考虑这种情况:

    t     : ---------1---------2---------3
    source: ------1---2------3---4----5--|
    window: ---------|---------|---------|
    spec  : ------1----------3-----------|
    enigma: ------1---2----------4-------|
    

    答案意味着您在第 1 项之后想要一秒钟什么都没有。之后的下一项是 3,然后一直保持沉默。这是测试代码:

    var scheduler = new TestScheduler();
    var source = scheduler.CreateColdObservable<int>(
        RxTest.OnNext(700.MsTicks(),  1),
        RxTest.OnNext(1100.MsTicks(), 2),
        RxTest.OnNext(1800.MsTicks(), 3),
        RxTest.OnNext(2200.MsTicks(), 4),
        RxTest.OnNext(2600.MsTicks(), 5),
        RxTest.OnCompleted<int>(3000.MsTicks())
    );
    
    var expectedResults = scheduler.CreateHotObservable<int>(
        RxTest.OnNext(700.MsTicks(),  1),
        RxTest.OnNext(1800.MsTicks(), 3),
        RxTest.OnCompleted<int>(3000.MsTicks())
    );
    
    var target = source
        .Window(() => Observable.Timer(TimeSpan.FromSeconds(1.0), scheduler))
        .SelectMany(xs => xs.Take(1));
    
    var observer = scheduler.CreateObserver<int>();
    target.Subscribe(observer);
    scheduler.Start();
    ReactiveAssert.AreElementsEqual(expectedResults.Messages, observer.Messages);
    

    我认为解决此问题的最佳方法是基于Scan 的解决方案,带有时间戳。您基本上将最后一条合法消息保存在内存中,并带有时间戳,如果新消息早一秒,则发出。否则,不要:

    public static IObservable<T> TrueThrottle<T>(this IObservable<T> source, TimeSpan span)
    {
        return TrueThrottle<T>(source, span, Scheduler.Default);
    }
    
    public static IObservable<T> TrueThrottle<T>(this IObservable<T> source, TimeSpan span, IScheduler scheduler)
    {
        return source
            .Timestamp(scheduler)
            .Scan(default(Timestamped<T>), (state, item) => state == default(Timestamped<T>) || item.Timestamp - state.Timestamp > span
                ? item
                : state
            )
            .DistinctUntilChanged()
            .Select(t => t.Value);
    }
    

    注意:测试代码使用 Nuget Microsoft.Reactive.Testing 和以下帮助类:

    public static class RxTest
    {
        public static long MsTicks(this int i)
        {
            return TimeSpan.FromMilliseconds(i).Ticks;
        }
    
        public static Recorded<Notification<T>> OnNext<T>(long msTicks, T t)
        {
            return new Recorded<Notification<T>>(msTicks, Notification.CreateOnNext(t));
        }
    
        public static Recorded<Notification<T>> OnCompleted<T>(long msTicks)
        {
            return new Recorded<Notification<T>>(msTicks, Notification.CreateOnCompleted<T>());
        }
    
        public static Recorded<Notification<T>> OnError<T>(long msTicks, Exception e)
        {
            return new Recorded<Notification<T>>(msTicks, Notification.CreateOnError<T>(e));
        }
    }
    

    【讨论】:

    • 我一定会去看看的。感谢您的精彩回答。
    猜你喜欢
    • 1970-01-01
    • 2017-12-13
    • 1970-01-01
    • 1970-01-01
    • 2013-03-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多