【发布时间】:2016-01-24 00:51:21
【问题描述】:
我是 Rx 的新手,所以我可能在这里犯了一些基本错误。
我想创建一个非常简单的套接字服务器,它可以使用 Observables 从客户端接收消息。为此我使用了 Rxx,它在 System.Net.Sockets 命名空间中提供了扩展方法,并且还提供了 ObserableTcpListener 静态工厂类。
这是我目前所拥有的,几乎是从各种来源偷来的:
IPEndPoint endpoint = new IPEndPoint(IPAddress.Parse("127.0.0.1"), 9001);
TcpListener listener = new TcpListener(endpoint);
IObservable<TcpClient> clients = listener
.StartSocketObservable(1)
.SelectMany<Socket, TcpClient>(socket => SocketToTcpClient(socket));
.Finally(listener.Stop)
clients.Subscribe(client =>
{
OnConnect(client).Subscribe(
message => OnMessage(client, message),
ex => OnException(client, ex),
() => OnCompleted(client));
});
private static IObservable<TcpClient> SocketToTcpClient(Socket socket)
{
TcpClient client = new TcpClient();
client.Client = socket;
return Observable.Return<TcpClient>(client);
}
private static IObservable<byte[]> OnConnect(TcpClient client)
{
return client.Client.ReceiveUntilCompleted(SocketFlags.None);
}
private static void OnMessage(TcpClient client, byte[] message)
{
Console.WriteLine("Mesage Received! - {0}", Encoding.UTF8.GetString(message));
}
private static void OnCompleted(TcpClient client)
{
Console.WriteLine("Completed.");
}
private static void OnException(TcpClient client, Exception ex)
{
Console.WriteLine("Exception: {0}", ex.ToString());
}
这行得通……在某种程度上。我可以建立一个客户端连接。一旦该连接终止,Observable 序列似乎就会终止并调用.Finally(listener.Stop)。显然,这不是我想要的。
我尝试使用 ObserableTcpListener.Start() 工厂类,但这让我得到了完全相同的结果。
IObservable<TcpClient> sockets = ObservableTcpListener.Start(endpoint);
sockets.Subscribe(client =>
{
OnConnect(client).Subscribe(
message => OnMessage(client, message),
ex => OnException(client, ex),
() => OnCompleted(client));
});
我想我确实理解这里的问题:clients 可观察序列在第一个客户端终止后只是空的,因此调用了 .Finally(listener.Stop)。
我需要做什么来规避这个问题?如何继续侦听传入连接?
【问题讨论】:
-
searchcode.com/codesearch/view/14317362 对我来说,代码看起来只接受一个连接。无论如何,建议放弃这种方法并使用标准技术来运行套接字。我看不出这比直接使用 TcpListener/Client 有任何优势,可能与 await 一起使用。
-
@usr 主要原因是,在大多数情况下,我确实喜欢以 Rx 方式编写代码,因为它可读性强,并且非常清楚地表达了一个人的意图。第二个原因是我正在编写一堆不同的事件样式代码,而 Rx 为所有这些样式提供了一个很好的抽象,能够保持模式相同。第三,我只是想学习 Rx。感谢您的建议!
-
如果我可以为此添加一个反论点:这种风格将正常的顺序代码拆分为回调,这对于代码质量来说通常是一件糟糕的事情。例如,您的代码中有一个经典错误,您假设 TCP 发送“消息”。它没有,因此如果你运气不好,
Encoding.UTF8.GetString有时会返回垃圾,并且你的“消息”在 UTF8 编码的代码点中途被分割。这种风格很难解决。在顺序代码中,您可以使用 StreamReader 或 BinaryReader 并提取数据。通过推动,您必须接受即将发生的事情。 -
看起来有人已经下定决心要如何编写代码。我想说这是函数响应式编程和 Rx 的一个很好的用例。继续前进:-)
-
@usr 我非常感谢您的关注。但是,请记住这不是实际的应用程序代码。这只是我试图让一些东西启动并运行。我很清楚没有“消息”,只有字节字符串(包含任何内容),并且不小心尝试将它们转换为 UTF8 字符串是非常危险的。这实际上是我有兴趣调查的问题之一(无论结果如何)。此外,我认为将 Rx 响应式风格与传统回调进行比较有点不公平。无论哪种方式,我都会记住您的 cmets。
标签: c# sockets system.reactive