【问题标题】:Observable TcpListener terminates after single connectionObservable TcpListener 在单次连接后终止
【发布时间】: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


【解决方案1】:

让您的Observable热门并在有订阅时坚持下去。

IObservable<TcpClient> clients = listener
    .StartSocketObservable(1)
    .SelectMany<Socket, TcpClient>(socket => SocketToTcpClient(socket))
    .Finally(listener.Stop)
    .Publish().RefCount();

【讨论】:

  • 很抱歉,我发了这个问题后病倒了一段时间,我完全忘记了我问过这个问题。不管怎样,我几个小时前就试过了,是的,这实际上似乎已经修复了它。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-08-30
  • 2014-11-03
  • 1970-01-01
  • 2012-01-16
  • 2015-01-05
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多