【问题标题】:Use of IObservable instead of events使用 IObservable 代替事件
【发布时间】:2012-06-25 13:04:43
【问题描述】:

我最近一直在阅读有关 IObservable 的文章。到目前为止,我已经查看了各种 SO 问题,并观看了有关他们可以做什么的视频。我认为的整个“推动”机制非常出色,但我仍在试图弄清楚一切究竟是做什么的。根据我的阅读,我猜在某种程度上 IObservable 是可以“观察”的东西,而 IObservers 是“观察者”。

所以现在我要尝试在我的应用程序中实现它。在开始之前,我想先弄清楚一些事情。我已经看到 IObservable 与 IEnumerable 相反,但是,在我的特定实例中,我真的看不到任何可以合并到我的应用程序中的地方。

目前,我大量使用事件,以至于我可以看到“管道”开始变得难以管理。我想,IObservable 可以帮助我。

考虑以下设计,它是我应用程序中 I/O 的包装器(仅供参考,我通常必须处理字符串):

我有一个名为 IDataIO 的基本接口:

public interface IDataIO
{
  event OnDataReceived;
  event OnTimeout:
  event OnTransmit;
}

现在,我目前有三个实现这个接口的类,每个类都在某种程度上利用了异步方法调用,引入了某种类型的多线程处理:

public class SerialIO : IDataIO;
public class UdpIO : IDataIO;
public class TcpIO : IDataIO;

这些类中的每一个都有一个实例封装到我的最终类中,称为 IO(它还实现了 IDataIO - 遵循我的策略模式):

public class IO : IDataIO
{
  public SerialIO Serial;
  public UdpIO Udp;
  public TcpIO Tcp;
}

我已经利用策略模式来封装这三个类,这样当在运行时在不同的IDataIO 实例之间进行更改时,它对最终用户来说是“不可见的”。正如您所想象的,这在后台导致了相当多的“事件管道”。

那么,在我的情况下,如何在此处使用“推送”通知?我不想订阅事件(DataReceived 等),而是简单地将数据推送给任何感兴趣的人。我有点不确定从哪里开始。我仍在尝试玩弄Subject 的想法/泛型类,以及它的各种化身(ReplaySubject/AsynSubject/BehaviourSubject)。有人可以请教我这个(也许参考我的设计)?或者这根本不适合IObservable

PS。请随时纠正我的任何“误解”:)

【问题讨论】:

    标签: c# events system.reactive


    【解决方案1】:

    Observables 非常适合表示数据流,因此您的 DataReceived 事件可以很好地模拟 observable 模式,例如 IObservable<byte>IObservable<byte[]>。您还可以获得OnErrorOnComplete 的额外好处,它们很方便。

    在实现方面,很难说你的具体场景,但我们经常使用Subject<T>作为底层源并调用OnNext来推送数据。也许像

    // Using a subject is probably the easiest way to push data to an Observable
    // It wraps up both IObservable and IObserver so you almost never use IObserver directly
    private readonly Subject<byte> subject = new Subject<byte>();
    
    private void OnPort_DataReceived(object sender, EventArgs e)
    {
        // This pushes the data to the IObserver, which is probably just a wrapper
        // around your subscribe delegate is you're using the Rx extensions
        this.subject.OnNext(port.Data); // pseudo code 
    }
    

    然后您可以通过属性公开主题:

    public IObservable<byte> DataObservable
    {
        get { return this.subject; } // Or this.subject.AsObservable();
    }
    

    您可以将IDataIO 上的DataReceived 事件替换为IObservable&lt;T&gt;,并让每个策略类以他们需要的任何方式处理其数据并将其推送到Subject&lt;T&gt;

    另一方面,订阅 Observable 的人可以像处理事件一样处理它(只需使用Action&lt;byte[]&gt;),或者您可以使用Select、@987654338 在流上执行一些非常有用的工作@、Buffer

    private IDataIO dataIo = new ...
    
    private void SubscribeToData()
    { 
        dataIo.DataObservable.Buffer(16).Subscribe(On16Bytes);
    }
    
    private void On16Bytes(IList<byte> bytes)
    {
        // do stuff
    }
    

    ReplaySubject/ConnectableObservables 当您知道您的订阅者将迟到参加聚会但仍需要赶上所有事件时,这非常棒。源缓存它推送的所有内容,并为每个订阅者重播所有内容。只有你可以说这是否是你真正需要的行为(但要小心,因为它会缓存所有显然会增加你的内存使用的东西)。

    当我学习 Rx 时,我发现关于 Rx 的 http://leecampbell.blogspot.co.uk/ 博客系列对于理解理论非常有用(这些帖子现在有点过时了,API 也发生了变化,所以请注意这一点)

    【讨论】:

    • 嗨 RichK,您能详细说明一下 Subject 属性吗?这是如何声明的?而这个类的用户,他们将如何“订阅” IObservable DataReceived
    • @Simon 我做了一些修改,如果您仍然不确定,请告诉我:)
    • 谢谢,这解决了一些问题。只有一件事,我假设 dataIo.DataObservablepublic IObservable&lt;byte&gt; DataObservable
    • 是的,抱歉,代码不太正确 - 我忘了命名该属性! [固定]
    • 听起来您需要了解一些有关可观察对象和反应模式的知识。反应是奇怪的。它是干什么用的?什么时候使用合适?我强烈推荐这个系列的文章。 introtorx.com/content/v1.0.10621.0/01_WhyRx.html#WhyRx
    【解决方案2】:

    这绝对是 observables 的理想案例。 IO 类可能会看到最大的改进。首先,让我们更改接口以使用可观察对象,看看组合类变得多么简单。

    public interface IDataIO
    {
        //you will have to fill in the types here.  Either the event args
        //the events provide now or byte[] or something relevant would be good.
        IObservable<???> DataReceived;
        IObservable<???> Timeout;
        IObservable<???> Transmit;
    }
    
    public class IO : IDataIO
    {
        public SerialIO Serial;
        public UdpIO Udp;
        public TcpIO Tcp;
    
        public IObservable<???> DataReceived
        {
            get 
            {
                return Observable.Merge(Serial.DataReceived,
                                        Udp.DataReceived,
                                        Tcp.DataReceived);
            }
        }
    
        //similarly for other two observables
    }
    

    旁注:你可能会注意到我更改了接口成员的名称。在 .NET 中,事件通常命名为 &lt;event name&gt;,引发它们的函数称为 On&lt;event name&gt;

    对于生产类,您有几个取决于实际来源的选项。假设您在SerialIO 中使用.NET SerialPort 类,并且DataReceived 返回IObservable&lt;byte[]&gt;。由于 SerialPort 已经有一个接收数据的事件,您可以直接使用它来制作您需要的 observable。

    public class SerialIO : IDataIO
    {
        private SerialPort _port;
    
        public IObservable<byte[]> DataRecived
        {
            get
            {
                return Observable.FromEventPattern<SerialDataReceivedEventHandler,
                                                   SerialDataReceivedEventArgs>(
                            h => _port.DataReceived += h,
                            h => _port.DataReceived -= h)
                       .Where(ep => ep.EventArgs.EventType == SerialData.Chars)
                       .Select(ep =>
                               {
                                  byte[] buffer = new byte[_port.BytesToRead];
                                  _port.Read(buffer, 0, buffer.Length);
                                  return buffer;
                               });
            }
        }
    }
    

    对于没有现有事件源的情况,您可能需要使用 RichK 建议的主题。他的回答很好地涵盖了这种使用模式,所以我不会在这里重复。

    您没有展示如何使用此接口,但根据用例,让这些类上的其他函数本身返回 IObservables 并完全取消这些“事件”可能更有意义。使用基于事件的异步模式,您必须将事件与调用的函数分开以触发工作,但使用可观察对象,您可以从函数中返回它们,以使您订阅的内容更加明显。该方法还允许从每个调用返回的 observables 发送 OnErrorOnCompleted 消息以表示操作结束。根据您对组合类的使用,我不希望这在这种特殊情况下有用,但请记住这一点。

    【讨论】:

    • +1 谢谢,这里有一些很好的信息。就在Merge() 语句上——这将一系列可观察对象合并为 1——在我的应用程序中,我只会一次使用其中一个(串行/udp/tcp)并允许用户在不同的界面之间切换(因此我对事件管道的困境)。是否建议在这里合并可观察对象?欣赏异步串行事件的链接:)
    • @Simon 或者您可以(最好)停止在其他 observables 上生成消息,以便 Merge 一次只接收一条消息。如果是这种情况,最好使用Switch 而不是Merge
    【解决方案3】:

    使用 IObservable 代替事件

    如果只对 nuget 包 rxx 具有的属性更改感兴趣:

    IObservable<string> obs=Observable2.FromPropertyChangedPattern(() => obj.Name)
    

    (以及许多其他方法)


    或者如果事件排除了属性更改/希望避免实施 INotifyPropertyChanged

    class ObserveEvent_Simple
    {
        public static event EventHandler SimpleEvent;
        static void Main()
        {          
           IObservable<string> eventAsObservable = Observable.FromEventPattern(
                ev => SimpleEvent += ev,
                ev => SimpleEvent -= ev);
        }
    }
    

    类似于 u/Gideon Engelberth 来自http://rxwiki.wikidot.com/101samples#toc6

    被覆盖 https://rehansaeed.com/reactive-extensions-part2-wrapping-events/


    这篇 codeproject 文章也致力于将事件转换为反应事件

    https://www.codeproject.com/Tips/1078183/Weak-events-in-NET-using-Reactive-Extensions-Rx

    并且还处理弱订阅

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多