【发布时间】: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<XxxEventHandler, XxxEventArgs>(h => xxx.Event += h, h => xxx.Event -= h)语法来观察事件。 -
很好的建议。切换到您建议的语法后,调用
WebSocket.MessageReceived += h时出现异常。显然Observable.FromEventPattern<T>(Object, string)吞下了异常。 -
我想吸取的教训是始终避免使用
Observable.FromEventPattern<T>(Object, string),而改用FromEventPattern<T>(Action<EventHandler<T>>, Action<EventHandler<T>>)之类的东西?异常吞咽行为浪费了我很多时间来找出问题所在。 -
我不知道它是否会吞下异常——它只是不是静态类型的,很容易出错。我总是使用
FromEventPattern<T>(Action<EventHandler<T>>, Action<EventHandler<T>>),因为它是静态检查的。 -
我将遵循同样的做法。感谢您的帮助!
标签: c# system.reactive