【问题标题】:Alternative to using Subject in reactive programming?在响应式编程中使用 Subject 的替代方案?
【发布时间】:2017-08-30 16:31:43
【问题描述】:

在响应式编程中,Subject 类型的使用通常是不受欢迎的。在以下情况下,我使用Subject 允许在创建通知的基础源之前订阅通知。有没有不使用Subject 的替代方法来实现这一点?

using System;
using System.Reactive;
using System.Reactive.Linq;
using System.Reactive.Subjects;
using System.Threading.Tasks;
using Windows.Networking.Sockets;
using Windows.Storage.Streams;

class Program
{
    public static void Main()
    {
        var socket = new ObservableMessageWebSocket();

        socket.Messages.Subscribe(Print); // Caller is allowed to subscribe before connect

        var uri = new Uri("ws://mydomain.com/messages");
        socket.ConnectAsync(uri).Wait(); // Caller is allowed to connect after subscribe

        Console.ReadLine();
    }

    public static void Print(string message)
    {
        Console.WriteLine(message);
    }
}

class ObservableMessageWebSocket
{
    // Is there a way to get rid of this Subject?
    private readonly Subject<string> subject = new Subject<string>();

    private MessageWebSocket webSocket;

    public IObservable<string> Messages => subject;

    public async Task ConnectAsync(Uri uri)
    {
        webSocket = new MessageWebSocket();

        webSocket.Control.MessageType = SocketMessageType.Utf8;

        Observable
            .FromEventPattern<MessageWebSocketMessageReceivedEventArgs>(webSocket, nameof(webSocket.MessageReceived))
            .Select(ReadString)
            .Subscribe(subject);

        await webSocket.ConnectAsync(uri);
    }

    private static string ReadString(EventPattern<MessageWebSocketMessageReceivedEventArgs> pattern)
    {
        using (var reader = pattern.EventArgs.GetDataReader())
        {
            reader.UnicodeEncoding = UnicodeEncoding.Utf8;

            return reader.ReadString(reader.UnconsumedBufferLength);
        }
    }
}

编辑:澄清一下,我有几个软件组件订阅ObservableMessageWebSocket.Messages 以获取推送通知。有些组件在调用ObservableMessageWebSocket.ConnectAsync 之前订阅,有些组件在调用之后订阅。

下面的代码避免了Subject,但不能正常运行。组件在连接后订阅,并且永远不会收到通知。

class ObservableMessageWebSocket
{
    private MessageWebSocket WebSocket { get; }

    public IObservable<string> Messages { get; }

    public ObservableMessageWebSocket()
    {
        WebSocket = new MessageWebSocket();
        WebSocket.Control.MessageType = SocketMessageType.Utf8;
        Messages = Observable
            .FromEventPattern<MessageWebSocketMessageReceivedEventArgs>(WebSocket, nameof(WebSocket.MessageReceived))
            .Select(ReadString);
    }

    private static string ReadString(EventPattern<MessageWebSocketMessageReceivedEventArgs> pattern)
    {
        using (var reader = pattern.EventArgs.GetDataReader())
        {
            reader.UnicodeEncoding = UnicodeEncoding.Utf8;

            return reader.ReadString(reader.UnconsumedBufferLength);
        }
    }

    public async Task ConnectAsync(Uri uri)
    {
        await WebSocket.ConnectAsync(uri);
    }
}

下面的代码也不起作用。相同的症状。

class ObservableMessageWebSocket
{
    private MessageWebSocket WebSocket { get; }

    public IObservable<string> Messages { get; }

    public ObservableMessageWebSocket()
    {
        WebSocket = new MessageWebSocket();
        WebSocket.Control.MessageType = SocketMessageType.Utf8;
        Messages = Observable.Create<string>(o => Observable
            .FromEventPattern<MessageWebSocketMessageReceivedEventArgs>(WebSocket, nameof(WebSocket.MessageReceived))
            .Select(ReadString)
            .Subscribe(o));
    }

    private static string ReadString(EventPattern<MessageWebSocketMessageReceivedEventArgs> pattern)
    {
        using (var reader = pattern.EventArgs.GetDataReader())
        {
            reader.UnicodeEncoding = UnicodeEncoding.Utf8;

            return reader.ReadString(reader.UnconsumedBufferLength);
        }
    }

    public async Task ConnectAsync(Uri uri)
    {
        await WebSocket.ConnectAsync(uri);
    }
}

不知何故,下面的代码有效。

class ObservableMessageWebSocket
{
    private MessageWebSocket WebSocket { get; }

    private event EventHandler<string> StringReceived;

    public IObservable<string> Messages { get; }

    public ObservableMessageWebSocket()
    {
        WebSocket = new MessageWebSocket();
        WebSocket.Control.MessageType = SocketMessageType.Utf8;
        WebSocket.MessageReceived += HandleEvent;
        Messages = Observable
            .FromEventPattern<string>(this, nameof(StringReceived))
            .Select(p => p.EventArgs);
    }

    private void HandleEvent(MessageWebSocket sender, MessageWebSocketMessageReceivedEventArgs args)
    {
        var handler = StringReceived;
        if (handler == null) return;
        string message;
        using (var reader = args.GetDataReader())
        {
            reader.UnicodeEncoding = UnicodeEncoding.Utf8;
            message= reader.ReadString(reader.UnconsumedBufferLength);
        }
        handler.Invoke(this, message);
    }

    public async Task ConnectAsync(Uri uri)
    {
        await WebSocket.ConnectAsync(uri);
    }
}

对我来说,这三个似乎都相似。为什么只有最后一个有效?

【问题讨论】:

  • 通过更新的代码,您确定您实际上订阅了测试代码中的 observable 吗?另外,尝试使用Observable.FromEventPattern&lt;XxxEventHandler, XxxEventArgs&gt;(h =&gt; xxx.Event += h, h =&gt; xxx.Event -= h) 语法来观察事件。
  • 很好的建议。切换到您建议的语法后,调用 WebSocket.MessageReceived += h 时出现异常。显然Observable.FromEventPattern&lt;T&gt;(Object, string) 吞下了异常。
  • 我想吸取的教训是始终避免使用Observable.FromEventPattern&lt;T&gt;(Object, string),而改用FromEventPattern&lt;T&gt;(Action&lt;EventHandler&lt;T&gt;&gt;, Action&lt;EventHandler&lt;T&gt;&gt;)之类的东西?异常吞咽行为浪费了我很多时间来找出问题所在。
  • 我不知道它是否会吞下异常——它只是不是静态类型的,很容易出错。我总是使用FromEventPattern&lt;T&gt;(Action&lt;EventHandler&lt;T&gt;&gt;, Action&lt;EventHandler&lt;T&gt;&gt;),因为它是静态检查的。
  • 我将遵循同样的做法。感谢您的帮助!

标签: c# system.reactive


【解决方案1】:

避免主题通常是个好主意。在您的代码中,您将主题直接暴露给调用代码。任何执行((Subject&lt;string&gt;)socket.Messages).OnCompleted(); 的消费者都会停止您的代码工作。

您还新建了一个 WebSocket,之后应该将其处理掉。

有一种方法可以驾驭主题并使其表现得更好。

试试这个:

public IObservable<string> Connect(Uri uri)
{
    return
        Observable
            .Using(
                () =>
                {
                    var webSocket = new MessageWebSocket();
                    webSocket.Control.MessageType = SocketMessageType.Utf8;
                    return webSocket;
                },
                webSocket =>
                    Observable
                        .FromAsync(() => webSocket.ConnectAsync(uri))
                        .SelectMany(u =>
                            Observable
                                .FromEventPattern<MessageWebSocketMessageReceivedEventArgs>(webSocket, nameof(webSocket.MessageReceived))
                                .SelectMany(pattern =>
                                    Observable
                                        .Using(
                                            () =>
                                            {
                                                var reader = pattern.EventArgs.GetDataReader();
                                                reader.UnicodeEncoding = UnicodeEncoding.UTF8;
                                                return reader;
                                            },
                                            reader => Observable.Return(reader.ReadString(reader.UnconsumedBufferLength))))));

}

以下是使用现有代码样式避免主题的方法:

public IObservable<string> ConnectAsync(Uri uri)
{
    return
        Observable
            .Create<string>(async o =>
            {
                var webSocket = new MessageWebSocket();

                webSocket.Control.MessageType = SocketMessageType.Utf8;

                var subscription = Observable
                    .FromEventPattern<MessageWebSocketMessageReceivedEventArgs>(webSocket, nameof(webSocket.MessageReceived))
                    .Select(ReadString)
                    .Subscribe(o);

                await webSocket.ConnectAsync(uri);

                return subscription;
            });
}

这是一个有效的快速测试:

void Main()
{
    Connect(new Uri("https://stackoverflow.com/")).Subscribe(x => Console.WriteLine(x.Substring(0, 24)));
}

public IObservable<string> Connect(Uri uri)
{
    return
        Observable
            .Create<string>(async o =>
            {
                var webClient = new WebClient();

                webClient.UseDefaultCredentials = true;

                var subscription =
                    Observable
                        .Using(
                            () => new CompositeDisposable(webClient, Disposable.Create(() => Console.WriteLine("Disposed!"))),
                            _ =>
                                Observable
                                    .FromEventPattern<DownloadStringCompletedEventHandler, DownloadStringCompletedEventArgs>(
                                        h => webClient.DownloadStringCompleted += h, h => webClient.DownloadStringCompleted -= h)
                                    .Take(1))
                    .Select(x => x.EventArgs.Result)
                    .Subscribe(o);

                await webClient.DownloadStringTaskAsync(uri);

                return subscription;
            });
}

请注意,"Disposed!" 的显示表明 WebClient 已被释放。

【讨论】:

  • 感谢您的回答。您对IDisposable 的评论是正确的;在生产代码中我确实实现了处置,但在 StackOverflow 上我删除了这些行以强调我的问题的要点,即如何在创建通知的基础源之前允许订阅通知。
  • @hwaien - 如果您想“在创建通知的基础源之前允许订阅通知”,那么您应该使用 Observable.Create - 这就是它的用途。然后它会避开Subject&lt;string&gt;。我稍后会尝试将其添加到我的答案中。
  • @hwaien - 我在我的代码中添加了一个更符合你想要的答案。
  • 感谢您的建议。我不知道Observable.Create。然而,在使用它之后,我发现它并没有给我我期望的结果。我更新了我的问题,包括我尝试使用Observable.Create
  • @hwaien - 我认为现在可能还有其他事情发生。你能看一下我对你的问题的评论吗?
【解决方案2】:

主题并不是普遍不好,而且我认为您使用它的方式没有什么大问题。我会参考Why are Subjects not recommended in .NET Reactive Extensions?RX Subjects - are they to be avoided? 对它们及其用途进行一些合理的讨论。

鉴于此,我建议如下(同时删除所有字段和公开的属性):

public async Task<IObservable<string>> ConnectAsync(Uri uri)
{
    webSocket = new MessageWebSocket();

    webSocket.Control.MessageType = SocketMessageType.Utf8;

    var toReturn = Observable
        .FromEventPattern<MessageWebSocketMessageReceivedEventArgs>(webSocket, nameof(webSocket.MessageReceived))
        .Select(ReadString);

    await webSocket.ConnectAsync(uri);
    return toReturn;
}

这样,如果有人调用ConnectAsync 两次,他们可以获得单独的 observables。

【讨论】:

  • 感谢您的回答。链接的讨论很有见地。您对多次致电ConnectAsync 的评论是正确的;在生产代码中,我确实有处理多个调用的逻辑,但在 StackOverflow 上,我删除了这些行以强调我的问题的要点。
  • 只是一个快速的,@Shlomo。此函数的返回类型是Task&lt;IObservable&lt;string&gt;&gt;,但您的代码似乎返回了Task&lt;IDisposable&gt;。我读错了吗?
  • 不,你没看错,我写错了。已编辑和更正。谢谢。
猜你喜欢
  • 2014-07-25
  • 2014-11-18
  • 1970-01-01
  • 2020-07-13
  • 2023-03-23
  • 2011-12-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多