【问题标题】:Implementing IObservable<T> from scratch从头开始实现 IObservable<T>
【发布时间】:2009-11-20 07:57:43
【问题描述】:

响应式扩展提供了许多帮助方法,用于将现有事件和异步操作转换为可观察对象,但是您将如何从头开始实现 IObservable

IEnumerable 具有可爱的 yield 关键字,使其实现起来非常简单。

实现 IObservable 的正确方法是什么?

我需要担心线程安全吗?

我知道支持在特定的同步上下文中回调,但这是我作为 IObservable 作者需要担心的事情还是以某种方式内置的?

更新:

这是我的 C# 版本的 Brian 的 F# 解决方案

using System;
using System.Linq;
using Microsoft.FSharp.Collections;

namespace Jesperll
{
    class Observable<T> : IObservable<T>, IDisposable where T : EventArgs
    {
        private FSharpMap<int, IObserver<T>> subscribers = 
                 FSharpMap<int, IObserver<T>>.Empty;
        private readonly object thisLock = new object();
        private int key;
        private bool isDisposed;

        public void Dispose()
        {
            Dispose(true);
        }

        protected virtual void Dispose(bool disposing)
        {
            if (disposing && !isDisposed)
            {
                OnCompleted();
                isDisposed = true;
            }
        }

        protected void OnNext(T value)
        {
            if (isDisposed)
            {
                throw new ObjectDisposedException("Observable<T>");
            }

            foreach (IObserver<T> observer in subscribers.Select(kv => kv.Value))
            {
                observer.OnNext(value);
            }
        }

        protected void OnError(Exception exception)
        {
            if (isDisposed)
            {
                throw new ObjectDisposedException("Observable<T>");
            }

            if (exception == null)
            {
                throw new ArgumentNullException("exception");
            }

            foreach (IObserver<T> observer in subscribers.Select(kv => kv.Value))
            {
                observer.OnError(exception);
            }
        }

        protected void OnCompleted()
        {
            if (isDisposed)
            {
                throw new ObjectDisposedException("Observable<T>");
            }

            foreach (IObserver<T> observer in subscribers.Select(kv => kv.Value))
            {
                observer.OnCompleted();
            }
        }

        public IDisposable Subscribe(IObserver<T> observer)
        {
            if (observer == null)
            {
                throw new ArgumentNullException("observer");
            }

            lock (thisLock)
            {
                int k = key++;
                subscribers = subscribers.Add(k, observer);
                return new AnonymousDisposable(() =>
                {
                    lock (thisLock)
                    {
                        subscribers = subscribers.Remove(k);
                    }
                });
            }
        }
    }

    class AnonymousDisposable : IDisposable
    {
        Action dispose;
        public AnonymousDisposable(Action dispose)
        {
            this.dispose = dispose;
        }

        public void Dispose()
        {
            dispose();
        }
    }
}

编辑:如果 Dispose 被调用两次,不要抛出 ObjectDisposedException

【问题讨论】:

  • Wes Dyer 现在在 Channel9 上有一个视频,讨论这些接口的合同。
  • 酷 - 一定会看的:)
  • 当 Dispose() 被调用两次时,你不应该抛出 ObjectDisposed。在调用 dispose 后调用其他方法时,您应该抛出 ObjectDisposed。
  • @JohnGietzen 你是对的。我编辑了代码以反映这一点。

标签: c# system.reactive


【解决方案1】:

official documentation 不赞成用户自己实现 IObservable。相反,用户应该使用工厂方法Observable.Create

如果可能,通过组合现有的运算符来实现新的运算符。否则使用 Observable.Create 实现自定义运算符

Observable.Create 恰好是 Reactive 内部类 AnonymousObservable 的一个简单包装器:

public static IObservable<TSource> Create<TSource>(Func<IObserver<TSource>, IDisposable> subscribe)
{
    if (subscribe == null)
    {
        throw new ArgumentNullException("subscribe");
    }
    return new AnonymousObservable<TSource>(subscribe);
}

我不知道他们为什么不公开他们的实施,但是,嘿,随便。

【讨论】:

  • 正确。不要自己实现IObservable&lt;T&gt;IObserver&lt;T&gt;
  • 嗨,李。喜欢你关于 RX 的书——以前的博客,比官方文档更好的指南。
  • 干杯。由于 Rx 现在已经开源,希望能够帮助团队更新官方代码/文档。
【解决方案2】:

老实说,我不确定这一切是否“正确”,但根据我目前的经验,感觉还不错。它是 F# 代码,但希望您能体会其中的味道。它允许您“新建”一个源对象,然后您可以调用 Next/Completed/Error on,它管理订阅并在源或客户端做坏事时尝试断言。

type ObservableSource<'T>() =     // '
    let protect f =
        let mutable ok = false
        try 
            f()
            ok <- true
        finally
            Debug.Assert(ok, "IObserver methods must not throw!")
            // TODO crash?
    let mutable key = 0
    // Why a Map and not a Dictionary?  Someone's OnNext() may unsubscribe, so we need threadsafe 'snapshots' of subscribers to Seq.iter over
    let mutable subscriptions = Map.empty : Map<int,IObserver<'T>>  // '
    let next(x) = subscriptions |> Seq.iter (fun (KeyValue(_,v)) -> protect (fun () -> v.OnNext(x)))
    let completed() = subscriptions |> Seq.iter (fun (KeyValue(_,v)) -> protect (fun () -> v.OnCompleted()))
    let error(e) = subscriptions |> Seq.iter (fun (KeyValue(_,v)) -> protect (fun () -> v.OnError(e)))
    let thisLock = new obj()
    let obs = 
        { new IObservable<'T> with       // '
            member this.Subscribe(o) =
                let k =
                    lock thisLock (fun () ->
                        let k = key
                        key <- key + 1
                        subscriptions <- subscriptions.Add(k, o)
                        k)
                { new IDisposable with 
                    member this.Dispose() = 
                        lock thisLock (fun () -> 
                            subscriptions <- subscriptions.Remove(k)) } }
    let mutable finished = false
    // The methods below are not thread-safe; the source ought not call these methods concurrently
    member this.Next(x) =
        Debug.Assert(not finished, "IObserver is already finished")
        next x
    member this.Completed() =
        Debug.Assert(not finished, "IObserver is already finished")
        finished <- true
        completed()
    member this.Error(e) =
        Debug.Assert(not finished, "IObserver is already finished")
        finished <- true
        error e
    // The object returned here is threadsafe; you can subscribe and unsubscribe (Dispose) concurrently from multiple threads
    member this.Value = obs

我会对任何关于这里好坏的想法感兴趣;我还没有机会看到来自 devlabs 的所有新的 Rx 东西......

我自己的经验表明:

  • 那些订阅 observables 的人永远不应该从订阅中抛出。当订阅者抛出异常时,observable 无法做任何合理的事情。 (这类似于事件。)异常很可能会冒泡到顶级的 catch-all 处理程序或使应用崩溃。
  • 源可能应该是“逻辑上单线程的”。我认为编写可以对并发 OnNext 调用做出反应的客户端可能更难;即使每个单独的调用来自不同的线程,也有助于避免并发调用。
  • 拥有一个执行某些“合同”的基类/助手类绝对有用。

我很好奇人们是否可以在这些方面提出更具体的建议。

【讨论】:

  • 谢谢,我在 C# 中创建了类似的东西,最终使用 F# Map 集合来避免枚举期间的锁定。另一种选择是使用 Eric Lippert 的 Immutable AVLTree 之类的东西。我已经说服自己,确保在适当的上下文中接收事件是观察者的责任,并且可观察者应该坚持每次都在同一个线程上引发事件(如您所写)。
【解决方案3】:

是的,yield 关键字很可爱;也许 IObservable(OfT) 会有类似的东西? [编辑:在 Eric Meijer 的 PDC '09 talk 中,他说“是的,注意这个空间”以声明性的收益来生成 observables。]

对于接近的东西(而不是自己滚动),请查看“(not yet) 101 Rx Samples”wiki 的the bottom,其中团队建议使用 Subject(T) 类作为“后端”来实现 IObservable(的)。这是他们的例子:

public class Order
{            
    private DateTime? _paidDate;

    private readonly Subject<Order> _paidSubj = new Subject<Order>();
    public IObservable<Order> Paid { get { return _paidSubj.AsObservable(); } }

    public void MarkPaid(DateTime paidDate)
    {
        _paidDate = paidDate;                
        _paidSubj.OnNext(this); // Raise PAID event
    }
}

private static void Main()
{
    var order = new Order();
    order.Paid.Subscribe(_ => Console.WriteLine("Paid")); // Subscribe

    order.MarkPaid(DateTime.Now);
}

【讨论】:

  • 恕我直言 Subject 绝对是您想要生成自己的 observable 的正确方法。
  • 顺便说一句,AsyncSubject 在这里是一个更好的选择,因为它为未来的订阅者保留了最后一个值。在您的示例中,必须在实际支付事件发生之前订阅。
  • @Nappy:我不知道AsyncSubject&lt;T&gt;——谢谢你提到它。
  • @Nappy 你的意思是 BehaviorSubject 而不是 AsyncSubject 吗?
【解决方案4】:
  1. 打开Reflector看看。

  2. 观看一些 C9 视频 - this 展示了如何“导出” Select 'combinator'

  3. 秘诀是创建 AnonymousObservable、AnonymousObserver 和 AnonymousDisposable 类(它们只是解决您无法实例化接口的问题)。当您使用 Actions 和 Funcs 传递时,它们包含零实现。

例如:

public class AnonymousObservable<T> : IObservable<T>
{
    private Func<IObserver<T>, IDisposable> _subscribe;
    public AnonymousObservable(Func<IObserver<T>, IDisposable> subscribe)
    {
        _subscribe = subscribe;
    }

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

剩下的交给你吧……这是一个很好的理解练习。

here 有一个不错的小线程,上面有相关问题。

【讨论】:

  • 谢谢,但不是很有帮助。我已经浏览过反射器和 C9 的大部分视频。 Reflector 只显示了实际的实现,很难从中推断出有关线程等的规则。此外,您所谓的秘密只是将正确实现的责任从实际的可观察类推到提供的 Func 上——它没有透露实现该 Func 的规则。所以基本上你什么也没告诉我,除了我自己弄清楚其余的 :)
  • 点了。老实说,到目前为止,我的大部分努力都在尝试编写他们所谓的“组合器”,而不是实际来源。您可以在此处从我的问题的答案中收集一些指导原则(目前获得“官方”答案的最佳地点):social.msdn.microsoft.com/Forums/en-US/rx/thread/…
【解决方案5】:

关于这个实现只有一点评论:

在 .net fw 4 中引入并发集合之后,使用 ConcurrentDictioary 而不是简单的字典可能会更好。

它节省了对集合的处理锁。

阿迪。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2013-03-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-12-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多