【问题标题】:Sending data over NetworkStream using multiple threads使用多个线程通过 NetworkStream 发送数据
【发布时间】:2016-08-01 19:31:26
【问题描述】:

我正在尝试构建一个命令行聊天室,服务器正在处理连接并将一个客户端的输入重复回所有其他客户端。 目前,服务器能够接收来自多个客户端的输入,但只能将信息单独发送回这些客户端。我认为我的问题是每个连接都在一个单独的线程上处理。我如何允许线程相互通信或能够向每个线程发送数据?

服务器代码:

namespace ConsoleApplication
{


    class TcpHelper
    {


        private static object _lock = new object();
        private static List<Task> _connections = new List<Task>();


        private static TcpListener listener { get; set; }
        private static bool accept { get; set; } = false;

        private static Task StartListener()
        {
            return Task.Run(async () =>
            {
                IPAddress address = IPAddress.Parse("127.0.0.1");
                int port = 5678;
                listener = new TcpListener(address, port);

                listener.Start();

                Console.WriteLine($"Server started. Listening to TCP clients at 127.0.0.1:{port}");

                while (true)
                {
                    var tcpClient = await listener.AcceptTcpClientAsync();
                    Console.WriteLine("Client has connected");
                    var task = StartHandleConnectionAsync(tcpClient);
                    if (task.IsFaulted)
                        task.Wait();
                }
            });
        }

        // Register and handle the connection
        private static async Task StartHandleConnectionAsync(TcpClient tcpClient)
        {
            // start the new connection task
            var connectionTask = HandleConnectionAsync(tcpClient);



            // add it to the list of pending task 
            lock (_lock)
                _connections.Add(connectionTask);

            // catch all errors of HandleConnectionAsync
            try
            {
                await connectionTask;

            }
            catch (Exception ex)
            {
                // log the error
                Console.WriteLine(ex.ToString());
            }
            finally
            {
                // remove pending task
                lock (_lock)
                    _connections.Remove(connectionTask);
            }
        }






        private static async Task HandleConnectionAsync(TcpClient client)
        {

            await Task.Yield();


            {
                using (var networkStream = client.GetStream())
                {

                    if (client != null)
                    {
                        Console.WriteLine("Client connected. Waiting for data.");



                        StreamReader streamreader = new StreamReader(networkStream);
                        StreamWriter streamwriter = new StreamWriter(networkStream);

                        string clientmessage = "";
                        string servermessage = "";


                        while (clientmessage != null && clientmessage != "quit")
                        {
                            clientmessage = await streamreader.ReadLineAsync();
                            Console.WriteLine(clientmessage);
                            servermessage = clientmessage;
                            streamwriter.WriteLine(servermessage);
                            streamwriter.Flush();


                        }
                        Console.WriteLine("Closing connection.");
                        networkStream.Dispose();
                    }
                }

            }

        }
        public static void Main(string[] args)
        {
            // Start the server 

            Console.WriteLine("Hit Ctrl-C to close the chat server");
            TcpHelper.StartListener().Wait();

        }

    }

}

客户代码:

namespace Client2
{
    public class Program
    {

        private static void clientConnect()
        {
            TcpClient socketForServer = new TcpClient();
            bool status = true;
            string userName;
            Console.Write("Input Username: ");
            userName = Console.ReadLine();

            try
            {
                IPAddress address = IPAddress.Parse("127.0.0.1");
                socketForServer.ConnectAsync(address, 5678);
                Console.WriteLine("Connected to Server");
            }
            catch
            {
                Console.WriteLine("Failed to Connect to server{0}:999", "localhost");
                return;
            }
            NetworkStream networkStream = socketForServer.GetStream();
            StreamReader streamreader = new StreamReader(networkStream);
            StreamWriter streamwriter = new StreamWriter(networkStream);
            try
            {
                string clientmessage = "";
                string servermessage = "";
                while (status)
                {
                    Console.Write(userName + ": ");
                    clientmessage = Console.ReadLine();
                    if ((clientmessage == "quit") || (clientmessage == "QUIT"))
                    {
                        status = false;
                        streamwriter.WriteLine("quit");
                        streamwriter.WriteLine(userName + " has left the conversation");
                        streamwriter.Flush();

                    }
                    if ((clientmessage != "quit") && (clientmessage != "quit"))
                    {
                        streamwriter.WriteLine(userName + ": " + clientmessage);
                        streamwriter.Flush();
                        servermessage = streamreader.ReadLine();
                        Console.WriteLine("Server:" + servermessage);
                    }
                }
            }
            catch
            {
                Console.WriteLine("Exception reading from the server");
            }
            streamreader.Dispose();
            networkStream.Dispose();
            streamwriter.Dispose();
        }
        public static void Main(string[] args)
        {
            clientConnect();
        }
    }
}

【问题讨论】:

    标签: c# multithreading async-await .net-core


    【解决方案1】:

    您的代码中的主要错误是您没有尝试将从一个客户端接收到的数据发送到其他连接的客户端。您的服务器中有_connections 列表,但列表中存储的唯一内容是用于连接的Task 对象,您甚至不使用这些对象。

    相反,您应该维护自己的连接列表,这样当您收到来自一个客户端的消息时,您可以将该消息重新传输给其他客户端。

    至少,这应该是List&lt;TcpClient&gt;,但因为您使用的是StreamReaderStreamWriter,所以您还需要初始化这些对象并将其存储在列表中。此外,您应该包括一个客户端标识符。一个明显的选择是客户端的名称(即用户输入的名称),但是您的示例在聊天协议中没有提供任何机制来传输该标识作为连接初始化的一部分,所以在我的示例(下)我只是使用一个简单的整数值。

    您发布的代码中还有一些其他不规范之处,例如:

    • 在一个全新的线程中启动一个任务,只是为了执行一些语句,使您能够启动一个异步操作。在我的示例中,我只是省略了代码的 Task.Run() 部分,因为它不是必需的。
    • 在为IsFaulted 返回时检查特定于连接的任务。由于在返回此Task 对象时实际上不可能发生任何 I/O,因此此逻辑几乎没有用处。对Wait() 的调用将引发异常,该异常将传播到主线程的Wait() 调用,从而终止服务器。但是如果出现任何其他错误,您不会终止服务器,因此不清楚您为什么要在此处执行此操作。
    • 有一个对Task.Yield() 的虚假呼叫。我不知道你想在那里完成什么,但不管它是什么,那句话都没有用。我只是删除了它。
    • 在您的客户端代码中,您仅在发送数据后才尝试从服务器接收数据。这是非常错误的;您希望客户在数据发送给他们后立即响应并接收数据。在我的版本中,我包含了一个简单的小匿名方法,该方法被立即调用以启动一个单独的消息接收循环,该循环将与主用户输入循环异步并发执行。
    • 同样在客户端代码中,您在“退出”消息之后发送“...已离开...”消息,该消息会导致服务器关闭连接。这意味着服务器永远不会真正收到“……已经离开……”消息。我颠倒了消息的顺序,以便“退出”始终是客户端发送的最后一件事。

    我的版本是这样的:

    服务器:

    class TcpHelper
    {
        class ClientData : IDisposable
        {
            private static int _nextId;
    
            public int ID { get; private set; }
            public TcpClient Client { get; private set; }
            public TextReader Reader { get; private set; }
            public TextWriter Writer { get; private set; }
    
            public ClientData(TcpClient client)
            {
                ID = _nextId++;
                Client = client;
    
                NetworkStream stream = client.GetStream();
    
                Reader = new StreamReader(stream);
                Writer = new StreamWriter(stream);
            }
    
            public void Dispose()
            {
                Writer.Close();
                Reader.Close();
                Client.Close();
            }
        }
    
        private static readonly object _lock = new object();
        private static readonly List<ClientData> _connections = new List<ClientData>();
    
        private static TcpListener listener { get; set; }
        private static bool accept { get; set; }
    
        public static async Task StartListener()
        {
            IPAddress address = IPAddress.Any;
            int port = 5678;
            listener = new TcpListener(address, port);
    
            listener.Start();
    
            Console.WriteLine("Server started. Listening to TCP clients on port {0}", port);
    
            while (true)
            {
                var tcpClient = await listener.AcceptTcpClientAsync();
                Console.WriteLine("Client has connected");
                var task = StartHandleConnectionAsync(tcpClient);
                if (task.IsFaulted)
                    task.Wait();
            }
        }
    
        // Register and handle the connection
        private static async Task StartHandleConnectionAsync(TcpClient tcpClient)
        {
            ClientData clientData = new ClientData(tcpClient);
    
            lock (_lock) _connections.Add(clientData);
    
            // catch all errors of HandleConnectionAsync
            try
            {
                await HandleConnectionAsync(clientData);
            }
            catch (Exception ex)
            {
                // log the error
                Console.WriteLine(ex.ToString());
            }
            finally
            {
                lock (_lock) _connections.Remove(clientData);
                clientData.Dispose();
            }
        }
    
        private static async Task HandleConnectionAsync(ClientData clientData)
        {
            Console.WriteLine("Client connected. Waiting for data.");
    
            string clientmessage;
    
            while ((clientmessage = await clientData.Reader.ReadLineAsync()) != null && clientmessage != "quit")
            {
                string message = "From " + clientData.ID + ": " + clientmessage;
    
                Console.WriteLine(message);
    
                lock (_lock)
                {
                    // Locking the entire operation ensures that a) none of the client objects
                    // are disposed before we can write to them, and b) all of the chat messages
                    // are received in the same order by all clients.
                    foreach (ClientData recipient in _connections.Where(r => r.ID != clientData.ID))
                    {
                        recipient.Writer.WriteLine(message);
                        recipient.Writer.Flush();
                    }
                }
            }
            Console.WriteLine("Closing connection.");
        }
    }
    

    客户:

    class Program
    {
        private const int _kport = 5678;
    
        private static async Task clientConnect()
        {
            IPAddress address = IPAddress.Loopback;
            TcpClient socketForServer = new TcpClient();
            string userName;
            Console.Write("Input Username: ");
            userName = Console.ReadLine();
    
            try
            {
                await socketForServer.ConnectAsync(address, _kport);
                Console.WriteLine("Connected to Server");
            }
            catch (Exception e)
            {
                Console.WriteLine("Failed to Connect to server {0}:{1}", address, _kport);
                return;
            }
    
    
            using (NetworkStream networkStream = socketForServer.GetStream())
            {
                var readTask = ((Func<Task>)(async () =>
                {
                    using (StreamReader reader = new StreamReader(networkStream))
                    {
                        string receivedText;
    
                        while ((receivedText = await reader.ReadLineAsync()) != null)
                        {
                            Console.WriteLine("Server:" + receivedText);
                        }
                    }
                }))();
    
                using (StreamWriter streamwriter = new StreamWriter(networkStream))
                {
                    try
                    {
                        while (true)
                        {
                            Console.Write(userName + ": ");
                            string clientmessage = Console.ReadLine();
                            if ((clientmessage == "quit") || (clientmessage == "QUIT"))
                            {
                                streamwriter.WriteLine(userName + " has left the conversation");
                                streamwriter.WriteLine("quit");
                                streamwriter.Flush();
                                break;
                            }
                            else
                            {
                                streamwriter.WriteLine(userName + ": " + clientmessage);
                                streamwriter.Flush();
                            }
                        }
    
                        await readTask;
                    }
                    catch (Exception e)
                    {
                        Console.WriteLine("Exception writing to server: " + e);
                        throw;
                    }
                }
            }
        }
    
        public static void Main(string[] args)
        {
            clientConnect().Wait();
        }
    }
    

    您还有很多工作要做。您可能希望在服务器端实现聊天用户名的正确初始化。至少,对于实际代码,您需要进行更多错误检查,并确保可靠地生成客户端 ID(如果您只想要正 ID 值,则不能超过 2^31-1回滚到0之前的连接)。

    我还进行了一些并非绝对必要的其他小改动,例如使用IPAddress.AnyIPAddress.Loopback 值而不是解析字符串,并且通常只是在这里和那里简化和清理代码。另外,我目前没有使用 C# 6 编译器,因此我更改了您使用 C# 6 功能的代码,以便改为使用 C# 5 进行编译。

    要做一个成熟的聊天服务器,你仍然需要完成你的工作。但我希望以上内容能让你重回正轨。

    【讨论】:

    • 哇。非常有帮助。是的,我主要是在谷歌上搜索和拼凑我发现的东西,并消除任何错误。你真的超越了我,我非常感激。我有点认为我需要存储客户,但我不确定这是否是多余的。非常感谢您解决这个问题!
    • @hereswilson:乐于提供帮助。请注意,以上只是您如何做事的一个示例。这不是您必须如何做的最后一句话。例如,您甚至可以选择将消息发送回发送它的客户端(稍微简化客户端枚举)。另请注意,我已经更改了中继部分的同步,同步整个操作,而不仅仅是让客户端发送到的部分。我决定我更喜欢这个(原因在代码中的注释中描述)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-10-24
    • 1970-01-01
    • 1970-01-01
    • 2017-12-03
    • 2015-10-19
    相关资源
    最近更新 更多