【问题标题】:How to buffer items when another observable emits true, and release them on false如何在另一个 observable 发出 true 时缓冲项目,并在 false 时释放它们
【发布时间】:2017-04-18 04:42:55
【问题描述】:

我有一个源流,通常希望在项目到达时发出它们。但是还有另一个可观察的——让我们称之为“门”。当门关闭时,源项应该缓冲,只有在门打开时才释放。

我已经能够编写一个函数来执行此操作,但它似乎比它需要的更复杂。我不得不使用Observable.Create 方法。我认为有一种方法可以使用Delay 或Buffer 方法使用几行更多功能代码来实现我的目标,但我不知道如何。 Delay 似乎特别有前途,但我不知道如何有时延迟,有时让一切立即通过(零延迟)。同样,我认为我可以使用Buffer,然后使用SelectMany;当门打开时,我会有长度为 1 的缓冲区,当门关闭时,我会有更长的缓冲区,但我还是不知道如何让它工作。

这是我构建的适用于所有测试的内容:

/// <summary>
/// Returns every item in <paramref name="source"/> in the order it was emitted, but starts
/// caching/buffering items when <paramref name="delay"/> emits true, and releases them when
/// <paramref name="delay"/> emits false.
/// </summary>
/// <param name="delay">
/// Functions as "gate" to start and stop the emitting of items. The gate is opened when true
/// and closed when false. The gate is open by default.
/// </param>

public static IObservable<T> DelayWhile<T>(this IObservable<T> source, IObservable<bool> delay) =>
    Observable.Create<T>(obs =>
    {
        ImmutableList<T> buffer = ImmutableList<T>.Empty;
        bool isDelayed = false;
        var conditionSubscription =
            delay
            .DistinctUntilChanged()
            .Subscribe(i =>
            {
                isDelayed = i;
                if (isDelayed == false)
                {
                    foreach (var j in buffer)
                    {
                        obs.OnNext(j);
                    }
                    buffer = ImmutableList<T>.Empty;
                }
            });
        var sourceSubscription =
            source
            .Subscribe(i =>
            {
                if (isDelayed)
                {
                    buffer = buffer.Add(i);
                }
                else
                {
                    obs.OnNext(i);
                }
            });
        return new CompositeDisposable(sourceSubscription, conditionSubscription);
    });

这是另一个通过测试的选项。它非常简洁,但不使用 Delay 或 Buffer 方法;我需要手动进行延迟/缓冲。

public static IObservable<T> DelayWhile<T>(this IObservable<T> source, IObservable<bool> delay) =>
    delay
    .StartWith(false)
    .DistinctUntilChanged()
    .CombineLatest(source, (d, i) => new { IsDelayed = d, Item = i })
    .Scan(
        seed: new { Items = ImmutableList<T>.Empty, IsDelayed = false },
        accumulator: (sum, next) => new
        {
            Items = (next.IsDelayed != sum.IsDelayed) ?
                    (next.IsDelayed ? sum.Items.Clear() : sum.Items) :
                    (sum.IsDelayed ? sum.Items.Add(next.Item) : sum.Items.Clear().Add(next.Item)),
            IsDelayed = next.IsDelayed
        })
    .Where(i => !i.IsDelayed)
    .SelectMany(i => i.Items);

这些是我的测试:

[DataTestMethod]
[DataRow("3-a 6-b 9-c", "1-f", "3-a 6-b 9-c", DisplayName = "Start with explicit no_delay, emit all future items")]
[DataRow("3-a 6-b 9-c", "1-f 2-f", "3-a 6-b 9-c", DisplayName = "Start with explicit no_delay+no_delay, emit all future items")]
[DataRow("3-a 6-b 9-c", "1-t", "", DisplayName = "Start with explicit delay, emit nothing")]
[DataRow("3-a 6-b 9-c", "1-t 2-t", "", DisplayName = "Start with explicit delay+delay, emit nothing")]
[DataRow("3-a 6-b 9-c", "5-t 10-f", "3-a 10-b 10-c", DisplayName = "When delay is removed, all cached items are emitted in order")]
[DataRow("3-a 6-b 9-c 12-d", "5-t 10-f", "3-a 10-b 10-c 12-d", DisplayName = "When delay is removed, all cached items are emitted in order")]
public void DelayWhile(string source, string isDelayed, string expectedOutput)
{
    (long time, string value) ParseEvent(string e)
    {
        var parts = e.Split('-');
        long time = long.Parse(parts[0]);
        string val = parts[1];
        return (time, val);
    }
    IEnumerable<(long time, string value)> ParseEvents(string s) => s.Split(new char[] { ' ' }, StringSplitOptions.RemoveEmptyEntries).Select(ParseEvent);
    var scheduler = new TestScheduler();
    var sourceEvents = ParseEvents(source).Select(i => OnNext(i.time, i.value)).ToArray();
    var sourceStream = scheduler.CreateHotObservable(sourceEvents);
    var isDelayedEvents = ParseEvents(isDelayed).Select(i => OnNext(i.time, i.value == "t")).ToArray();
    var isDelayedStream = scheduler.CreateHotObservable(isDelayedEvents);
    var expected = ParseEvents(expectedOutput).Select(i => OnNext(i.time, i.value)).ToArray();
    var obs = scheduler.CreateObserver<string>();
    var result = sourceStream.DelayWhile(isDelayedStream);
    result.Subscribe(obs);
    scheduler.AdvanceTo(long.MaxValue);
    ReactiveAssert.AreElementsEqual(expected, obs.Messages);
}

[TestMethod]
public void DelayWhile_SubscribeToSourceObservablesOnlyOnce()
{
    var scheduler = new TestScheduler();
    var source = scheduler.CreateHotObservable<int>();
    var delay = scheduler.CreateHotObservable<bool>();

    // No subscriptions until subscribe
    var result = source.DelayWhile(delay);
    Assert.AreEqual(0, source.ActiveSubscriptions());
    Assert.AreEqual(0, delay.ActiveSubscriptions());

    // Subscribe once to each
    var obs = scheduler.CreateObserver<int>();
    var sub = result.Subscribe(obs);
    Assert.AreEqual(1, source.ActiveSubscriptions());
    Assert.AreEqual(1, delay.ActiveSubscriptions());

    // Dispose subscriptions when subscription is disposed
    sub.Dispose();
    Assert.AreEqual(0, source.ActiveSubscriptions());
    Assert.AreEqual(0, delay.ActiveSubscriptions());
}

[TestMethod]
public void DelayWhile_WhenSubscribeWithNoDelay_EmitCurrentValue()
{
    var source = new BehaviorSubject<int>(1);
    var emittedValues = new List<int>();
    source.DelayWhile(Observable.Return(false)).Subscribe(i => emittedValues.Add(i));
    Assert.AreEqual(1, emittedValues.Single());
}

// Subscription timing issue?
[TestMethod]
public void DelayWhile_WhenSubscribeWithDelay_EmitNothing()
{
    var source = new BehaviorSubject<int>(1);
    var emittedValues = new List<int>();
    source.DelayWhile(Observable.Return(true)).Subscribe(i => emittedValues.Add(i));
    Assert.AreEqual(0, emittedValues.Count);
}

[TestMethod]
public void DelayWhile_CoreScenario()
{
    var source = new BehaviorSubject<int>(1);
    var delay = new BehaviorSubject<bool>(false);
    var emittedValues = new List<int>();

    // Since no delay when subscribing, emit value
    source.DelayWhile(delay).Subscribe(i => emittedValues.Add(i));
    Assert.AreEqual(1, emittedValues.Single());

    // Turn on delay and buffer up a few; nothing emitted
    delay.OnNext(true);
    source.OnNext(2);
    source.OnNext(3);
    Assert.AreEqual(1, emittedValues.Single());

    // Turn off delay; should release the buffered items
    delay.OnNext(false);
    Assert.IsTrue(emittedValues.SequenceEqual(new int[] { 1, 2, 3 }));
}

【问题讨论】:

    标签: system.reactive


    【解决方案1】:

    编辑:我忘记了在使用基于Join 和Join 的运算符(例如WithLatestFrom)时,当有两个冷可观察对象时会遇到的问题。不用说,下面提到的关于缺乏交易的批评比以往任何时候都更加明显。

    我会推荐这个,它更像我原来的解决方案,但使用了Delay 重载。它通过了除DelayWhile_WhenSubscribeWithDelay_EmitNothing 之外的所有测试。为了解决这个问题,我将创建一个接受初始默认值的重载:

    public static IObservable<T> DelayWhile<T>(this IObservable<T> source, IObservable<bool> delay, bool isGateClosedToStart)
    {
        return source.Publish(_source => delay
            .DistinctUntilChanged()
            .StartWith(isGateClosedToStart)
            .Publish(_delay => _delay
                .Select(isGateClosed => isGateClosed
                    ? _source.TakeUntil(_delay).Delay(_ => _delay)
                    : _source.TakeUntil(_delay)
                )
                .Merge()
            )
        );
    }
    
    public static IObservable<T> DelayWhile<T>(this IObservable<T> source, IObservable<bool> delay)
    {
        return DelayWhile(source, delay, false);
    }
    

    旧答案:

    我最近读了一本书,批评 Rx 不支持交易,我第一次尝试解决这个问题就是一个很好的例子:

    public static IObservable<T> DelayWhile<T>(this IObservable<T> source, IObservable<bool> delay)
    {
        return source.Publish(_source => delay
            .DistinctUntilChanged()
            .StartWith(false)
            .Publish(_delay => _delay
                .Select(isGateClosed => isGateClosed 
                    ? _source.Buffer(_delay).SelectMany(l => l) 
                    : _source)
                .Switch()
            )
        );
    }
    

    应该工作,除了有太多东西依赖于 delay observable,订阅顺序很重要:在这种情况下,Switch 在 Buffer 结束之前切换,所以什么都没有当延迟门关闭时最终会出来。

    这可以修复如下:

    public static IObservable<T> DelayWhile<T>(this IObservable<T> source, IObservable<bool> delay)
    {
        return source.Publish(_source => delay
            .DistinctUntilChanged()
            .StartWith(false)
            .Publish(_delay => _delay
                .Select(isGateClosed => isGateClosed 
                    ? _source.TakeUntil(_delay).Buffer(_delay).SelectMany(l => l) 
                    : _source.TakeUntil(_delay)
                )
                .Merge()
            )
        );
    }
    

    我的下一次尝试通过了你所有的测试,并且也使用了你想要的 Observable.Delay 重载:

    public static IObservable<T> DelayWhile<T>(this IObservable<T> source, IObservable<bool> delay)
    {
        return delay
            .DistinctUntilChanged()
            .StartWith(false)
            .Publish(_delay => source
                .Join(_delay,
                    s => Observable.Empty<Unit>(),
                    d => _delay,
                    (item, isGateClosed) => isGateClosed 
                        ? Observable.Return(item).Delay(, _ => _delay) 
                        : Observable.Return(item)
                )
                .Merge()
        );
    }
    

    Join 可以像这样简化为 WithLatestFrom:

    public static IObservable<T> DelayWhile<T>(this IObservable<T> source, IObservable<bool> delay)
    {
        return delay
            .DistinctUntilChanged()
            .StartWith(false)
            .Publish(_delay => source
                .WithLatestFrom(_delay,
                    (item, isGateClosed) => isGateClosed 
                        ? Observable.Return(item).Delay(_ => _delay) 
                        : Observable.Return(item)
                )
                .Merge()
        );
    }
    

    【讨论】:

    • 我最喜欢的是最后一个——我能理解。那个“递归”的东西——在传递给 Publish 的函数中使用了 _delay——是我在尝试使用 Delay 方法时遇到的问题。我不知道您可以将函数传递给 Publish,这非常有用。我没有意识到有一个 Publish 版本返回 IObservable 而不是 IConnectableObservable。
    • 嗯...您的解决方案不适用于非常简单的测试用例。我的前几个解决方案确实有效。 [TestMethod] public void DelayWhile_SimpleTest() { var source = new BehaviorSubject&lt;int&gt;(1); var emittedValues = new List&lt;int&gt;(); source.DelayWhile(Observable.Return(false)).Subscribe(i =&gt; emittedValues.Add(i)); Assert.AreEqual(1, emittedValues.Single()); }
    • 我在原帖中添加了更多测试用例。
    • 更新答案。
    • 感谢您的建议,但我想我会坚持使用我原来的 Observable.Create 代码。从外观上可以清楚地看出它会起作用。我不完全理解为什么其他选项(除了你的最后一个建议)不能很好地工作是挑剔或挂起。也许是您正在谈论的那些时间问题或交易问题。您的最后一个建议非常有效,但看起来有点复杂。
    【解决方案2】:

    建议的简洁答案。看起来它应该可以工作,但它没有通过所有测试。

    public static IObservable<T> DelayWhile<T>(this IObservable<T> source, IObservable<bool> delay)
    {
        source = source.Publish().RefCount();
        delay = delay.Publish().RefCount();
        var delayRemoved = delay.Where(i => i == false);
        var sourceWhenNoDelay = source.WithLatestFrom(delay.StartWith(false), (s, d) => d).Where(i => !i);
        return
            source
            .Buffer(bufferClosingSelector: () => delayRemoved.Merge(sourceWhenNoDelay))
            .SelectMany(i => i);
    }
    

    【讨论】:

    • 其实这行不通;运行测试时似乎挂起。所以现在唯一可行的解​​决方案是我原来的 LONG 与 Observable.Create 的解决方案。
    猜你喜欢
    • 1970-01-01
    • 2014-11-03
    • 2018-06-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多