【问题标题】:Rx.NET "gate" operatorRx.NET“门”操作员
【发布时间】:2018-06-05 07:58:06
【问题描述】:

[注意:如果这很重要,我使用的是 3.1。另外,我在 codereview 上问过这个问题,但到目前为止没有任何回应。]

我需要一个运算符来允许布尔值流充当另一个流的门(当门流为真时让值通过,当门流为假时丢弃它们)。我通常会为此使用 Switch,但如果源流很冷,它会继续重新创建它,这是我不想要的。

我也想自己清理一下,以便在源或门完成时结果完成。

public static IObservable<T> When<T>(this IObservable<T> source, IObservable<bool> gate)
{
    var s = source.Publish().RefCount();
    var g = gate.Publish().RefCount();

    var sourceCompleted = s.TakeLast(1).DefaultIfEmpty().Select(_ => Unit.Default);
    var gateCompleted = g.TakeLast(1).DefaultIfEmpty().Select(_ => Unit.Default);

    var anyCompleted = Observable.Amb(sourceCompleted, gateCompleted);

    var flag = false;
    g.TakeUntil(anyCompleted).Subscribe(value => flag = value);

    return s.Where(_ => flag).TakeUntil(anyCompleted);
}

除了整体冗长之外,我不喜欢订阅门,即使结果从未订阅过(在这种情况下,此运算符应该是空操作)。有没有办法摆脱那个订阅?

我也尝试过这种实现,但在自行清理时更糟糕:

return Observable.Create<T>(
    o =>
    {
        var flag = false;
        gate.Subscribe(value => flag = value);

        return source.Subscribe(
            value =>
            {
                if (flag) o.OnNext(value);
            });
    });

这些是我用来检查实现的测试:

[TestMethod]
public void TestMethod1()
{
    var output = new List<int>();

    var source = new Subject<int>();
    var gate = new Subject<bool>();

    var result = source.When(gate);
    result.Subscribe(output.Add, () => output.Add(-1));

    // the gate starts with false, so the source events are ignored
    source.OnNext(1);
    source.OnNext(2);
    source.OnNext(3);
    CollectionAssert.AreEqual(new int[0], output);

    // setting the gate to true will let the source events pass
    gate.OnNext(true);
    source.OnNext(4);
    CollectionAssert.AreEqual(new[] { 4 }, output);
    source.OnNext(5);
    CollectionAssert.AreEqual(new[] { 4, 5 }, output);

    // setting the gate to false stops source events from propagating again
    gate.OnNext(false);
    source.OnNext(6);
    source.OnNext(7);
    CollectionAssert.AreEqual(new[] { 4, 5 }, output);

    // completing the source also completes the result
    source.OnCompleted();
    CollectionAssert.AreEqual(new[] { 4, 5, -1 }, output);
}

[TestMethod]
public void TestMethod2()
{
    // completing the gate also completes the result
    var output = new List<int>();

    var source = new Subject<int>();
    var gate = new Subject<bool>();

    var result = source.When(gate);
    result.Subscribe(output.Add, () => output.Add(-1));

    gate.OnCompleted();
    CollectionAssert.AreEqual(new[] { -1 }, output);
}

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    更新:当门终止时也终止。我在复制/粘贴中错过了TestMethod2

        return gate.Publish(_gate => source
            .WithLatestFrom(_gate.StartWith(false), (value, b) => (value, b))
            .Where(t => t.b)
            .Select(t => t.value)
            .TakeUntil(_gate.IgnoreElements().Materialize()
        ));
    

    这通过了你的测试 TestMethod1,它不会在可观察门完成时终止。

    public static IObservable<T> When<T>(this IObservable<T> source, IObservable<bool> gate)
    {
        return source
            .WithLatestFrom(gate.StartWith(false), (value, b) => (value, b))
            .Where(t => t.b)
            .Select(t => t.value);
    }
    

    【讨论】:

    • 谢谢;这是非常优雅的(我喜欢),但它实际上没有通过第二次测试(当门通过时结果没有完成)。你看到它通过了吗?
    • 抱歉,我没看到TestMethod2。修改后的答案。
    • WithLatestFrom 已经是System.Reactive 某些版本的一部分
    【解决方案2】:

    这行得通:

    public static IObservable<T> When<T>(this IObservable<T> source, IObservable<bool> gate)
    {
        return
            source.Publish(ss => gate.Publish(gs =>
                gs
                    .Select(g => g ? ss : ss.IgnoreElements())
                    .Switch()
                    .TakeUntil(Observable.Amb(
                        ss.Select(s => true).Materialize().LastAsync(),
                        gs.Materialize().LastAsync()))));
    }
    

    这两项测试都通过了。

    【讨论】:

    • 另一个优雅的回应。你们是怎么知道所有这些东西的,我想我已经完成了两打 Rx 课程,但我仍然不知道很多:)
    • 如果您将 Publish 中的 lambda 更改为此,它可以工作:var gg = gate.Publish().RefCount(); var bothCompleted = Observable.Amb(ss.WhenCompleted(), gg.WhenCompleted()); return gate.Select(g =&gt; g ? ss : ss.IgnoreElements()).Switch().TakeUntil(bothCompleted); 其中WhenCompleted 只是.Select(_ =&gt; Unit.Default).IgnoreElements().Concat(Observable.Return(Unit.Default))
    • @MarcelPopescu - 当心.Publish().RefCount() - 它会创建脆弱的、运行一次的、可观察的。他们可以轻松地制作一个在第二次失败后看起来运行良好的 observable。通常最好封装在 .Publish(inner =&gt; { }) 中。
    • @MarcelPopescu - 我采用了您的查询的变体并使其正常工作。
    • 感谢您的警告;我之前确实遇到过.Publish().RefCount() 的问题。我会尽量避免。
    【解决方案3】:

    您与Observable.Create 走在正确的轨道上。您应该从 observable 的两个订阅中调用 onError 和 onCompleted 以在需要时正确完成或出错。此外,如果您打算在 sourcegate 完成之前处置 When 订阅,则通过在 Create 委托中返回两个 IDisposables,您可以确保正确清理两个订阅。

        public static IObservable<T> When<T>(this IObservable<T> source, IObservable<bool> gate)
        {
            return Observable.Create<T>(
                o =>
                {
                    var flag = false;
                    var gs = gate.Subscribe(
                        value => flag = value,
                        e => o.OnError(e),
                        () => o.OnCompleted());
    
                    var ss = source.Subscribe(
                        value =>
                        {
                            if (flag) o.OnNext(value);
                        },
                        e => o.OnError(e), 
                        () => o.OnCompleted());
    
                    return new CompositeDisposable(gs, ss);
                });
        }
    

    仅使用 Rx 运算符的更短但更难阅读的版本。对于冷的 observable,它可能需要源的发布/引用计数。

        public static IObservable<T> When<T>(this IObservable<T> source, IObservable<bool> gate)
        {
            return gate
                .Select(g => g ? source : source.IgnoreElements())
                .Switch()
                .TakeUntil(source.Materialize()
                                 .Where(s => s.Kind == NotificationKind.OnCompleted));
        }
    

    【讨论】:

    • 谢谢;这是完美的。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2010-11-25
    • 2015-07-25
    • 2011-05-10
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多