【问题标题】:How to implement my own operator in rx.net如何在 rx.net 中实现我自己的操作符
【发布时间】:2019-11-18 23:37:52
【问题描述】:

我需要 RX 中的磁滞滤波器功能。只有当先前发出的值和当前输入值相差一定量时,它才应该从源流中发出一个值。作为一个通用的扩展方法,它可以有如下签名:

public static IObservable<T> HysteresisFilter<T>(this IObservable<t> source, Func<T/*previously emitted*/, T/*current*/, bool> filter)

我无法弄清楚如何使用现有的运营商来实现这一点。我一直在寻找来自 RxJava 的 lift 之类的东西,以及任何其他创建我自己的运算符的方法。我看过这个checklist,但是我在网上没有找到任何例子。

以下方法(实际上都是相同的)对我来说似乎可以解决问题,但是是否有更多的 Rx 方法 可以做到这一点,例如不包装 subject 或实际实现运算符?

async Task Main()
{
    var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5));

    var rnd = new Random();
    var s = Observable.Interval(TimeSpan.FromMilliseconds(10))
            .Scan(0d, (a,_) => a + rnd.NextDouble() - 0.5)
            .Publish()
            .AutoConnect()
            ;

    s.Subscribe(Console.WriteLine, cts.Token);

    s.HysteresisFilter((p, c) => Math.Abs(p - c) > 1d).Subscribe(x => Console.WriteLine($"1> {x}"), cts.Token);
    s.HysteresisFilter2((p, c) => Math.Abs(p - c) > 1d).Subscribe(x => Console.WriteLine($"2> {x}"), cts.Token);

    await Task.Delay(Timeout.InfiniteTimeSpan, cts.Token).ContinueWith(_=>_, TaskContinuationOptions.OnlyOnCanceled);
}

public static class ReactiveOperators
{
    public static IObservable<T> HysteresisFilter<T>(this IObservable<T> source, Func<T, T, bool> filter)
    {
        return new InternalHysteresisFilter<T>(source, filter).AsObservable; 
    }

    public static IObservable<T> HysteresisFilter2<T>(this IObservable<T> source, Func<T, T, bool> filter)
    {
        var subject = new Subject<T>();
        T lastEmitted = default;
        bool emitted = false;

        source.Subscribe(
            value =>
            {
                if (!emitted || filter(lastEmitted, value))
                {
                    subject.OnNext(value);
                    lastEmitted = value;
                    emitted = true;
                }
            } 
            , ex => subject.OnError(ex)
            , () => subject.OnCompleted()
        );

        return subject;
    }

    private class InternalHysteresisFilter<T>: IObserver<T>
    {
        Func<T, T, bool> filter;
        T lastEmitted;
        bool emitted;

        private readonly Subject<T> subject = new Subject<T>();

        public IObservable<T> AsObservable => subject;

        public InternalHysteresisFilter(IObservable<T> source, Func<T, T, bool> filter)
        {
            this.filter = filter;
            source.Subscribe(this);
        }

        public IDisposable Subscribe(IObserver<T> observer)
        {
            return subject.Subscribe(observer);
        }

        public void OnNext(T value)
        {
            if (!emitted || filter(lastEmitted, value))
            {
                subject.OnNext(value);
                lastEmitted = value;
                emitted = true;
            }
        }

        public void OnError(Exception error)
        {
            subject.OnError(error);
        }

        public void OnCompleted()
        {
            subject.OnCompleted();
        }
    }
}

旁注:将有数千个此类过滤器应用于尽可能多的流。我需要吞吐量超过延迟,因此我正在寻找 CPU 和内存中开销最小的解决方案,即使其他人看起来更漂亮。

【问题讨论】:

    标签: c# system.reactive rx.net


    【解决方案1】:

    我在Introduction to Rx 一书中看到的大多数示例都使用Observable.Create 方法来创建新的运算符。

    Create 工厂方法是实现自定义可观察序列的首选方法。主题的使用应主要停留在样本和测试领域。 (citation)

    public static IObservable<T> HysteresisFilter<T>(this IObservable<T> source,
        Func<T, T, bool> predicate)
    {
        return Observable.Create<T>(observer =>
        {
            T lastEmitted = default;
            bool emitted = false;
            return source.Subscribe(value =>
            {
                if (!emitted || predicate(lastEmitted, value))
                {
                    observer.OnNext(value);
                    lastEmitted = value;
                    emitted = true;
                }
            }, observer.OnError, observer.OnCompleted);
        });
    }
    

    【讨论】:

    • 是的,这听起来很合理。这也意味着Rx.Net中没有用户定义的运算符这样的东西?但是,这种方法看起来既干净又便宜。
    • 我之前没有听说过“用户定义的运算符”这个词,所以它可能不是一个既定的术语。但不要相信我的话,因为我不是经验丰富的 RX 用户。 :-)
    • 可能不是既定术语,但我正在寻找类似于我在第一个链接文章中引用的内容。为此,它看起来是 RxJava 中的 Operator 接口——不过我自己从未使用过。
    • 在对象浏览器 (Ctrl+W,J) 中搜索“操作员”没有发现任何与 RX 相关的内容,所以我猜想这个概念还没有被 .NET 采用。除非它包含在单独的包中,但我在 NuGet 包管理器中也看不到任何相关内容。
    • Operator 不是 Rx.net 中的接口。您可以粗略地认为它是任何接受IObservable 并返回IObservable 的函数。
    【解决方案2】:

    这个答案与@Theodor 的答案相同,但它避免使用我通常会避免的Observable.Create

    public static IObservable<T> HysteresisFilter2<T>(this IObservable<T> source,
        Func<T, T, bool> predicate)
    {
        return source
            .Scan((emitted: default(T), isFirstItem: true, emit: false), (state, newItem) => state.isFirstItem || predicate(state.emitted, newItem)
                ? (newItem, false, true)
                : (state.emitted, false, false)
            )
            .Where(t => t.emit)
            .Select(t => t.emitted);
    }
    

    .Scan 是您在跟踪 observable 中项目的状态时想要使用的。

    【讨论】:

    • 嗨什洛莫。回避Observable.Create的原因是什么?是否对线程安全或此方法的性能有任何担忧?
    • 感谢您的建议。但我实际上对这种方法存在性能问题,即使我想拥有像这个一样的简洁的功能解决方案。创建 valuetuples 很便宜,但仍然有它的价格,我们在这里连续三个运算符。我不确定优化器在这里有什么机会。不幸的是,我必须考虑任何少量的 CPU 和内存。我将进行比较基准测试,看看这两种方法的表现如何。
    • @TheodorZoulias,当将可变编程与 System.Reactive 之类的函数库混合时,通常很容易搞砸可变编程。例如,如果您在 Observable.Create 之外声明了两个字段变量,那么您将遇到多重订阅问题。
    • @ZorgoZ,我认为另一个答案表现更好。如果您在计算毫秒数,那么您可能使用了错误的库和错误的语言。
    • 介于两者之间。据我所知,Rx 非常有效。但是我不是在制作低级协议之类的,没有 RT,因此不需要极高的精度或计时,但是我可以在某个硬件上处理的流的数量将很大程度上取决于我通过保持性能之间的平衡使用了多少资源和灵活性。由于延迟是可以接受的,我一般不关心 CLR。​​
    猜你喜欢
    • 2020-09-06
    • 1970-01-01
    • 2016-01-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-09-08
    • 1970-01-01
    相关资源
    最近更新 更多