如何在C#中实现类似C语言select的多IO并发聊天客户端?
C#实现同时监听控制台输入与Socket服务器消息的方案
需求与C语言参考实现
我需要编写一个聊天服务器的简单客户端,要求能同时监听控制台输入与Socket连接的服务器消息。以下是C语言中使用select实现该逻辑的示例,希望在C#中复刻该行为:
while (true) { //make set of descriptors fd_set reads; FD_ZERO(&reads); FD_SET(client_socket, &reads); FD_SET(0, &reads); select(client_socket+1, &reads, NULL, NULL, NULL); //handle user input from stdin if (FD_ISSET(0, &reads)) { /// } //handle message from server if (FD_ISSET(client_socket, &reads)) { /// } }
自行实现的Task.WhenAny方案
我自己编写了基于Task.WhenAny的C#实现,但不确定其有效性与效率,请问还有其他实现方式吗?
var ongoingTasks = new List<Task> {getServerMsg, getUserMsg}; while (true) { var finishedTask = await Task.WhenAny(ongoingTasks); if (finishedTask == getServerMsg) { OutputServerMsg(); ongoingTasks.Remove(getServerMsg); getServerMsg = ReceiveMessageFromServer(); ongoingTasks.Add(getServerMsg); } else if (finishedTask == getUserInput) { SendMsgToServer(); ongoingTasks.Remove(getUserInput); getUserInput = ReceiveInputFromStdIn(); ongoingTasks.Add(getUserInput); } }
你的Task.WhenAny方案评价
这个实现逻辑是有效的,完全能满足聊天客户端的需求:
- 核心思路通过
Task.WhenAny等待任一操作完成,符合异步编程的最佳实践 - 对于聊天客户端这种低并发场景,效率完全够用,
Task.WhenAny的开销可以忽略不计 - 可优化点:无需用
List<Task>维护任务,直接持有两个任务变量即可;建议添加异常处理逻辑,避免单个任务抛出异常导致循环终止
其他可行的C#实现方式
1. 使用Socket.Select(最贴近C语言select的实现)
C#的Socket类原生提供了Select静态方法,和C语言的select逻辑几乎一致,适合习惯原生IO多路复用的场景:
using System; using System.Net.Sockets; using System.Text; class ChatClient { static void Main(string[] args) { // 连接服务器 using var clientSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); clientSocket.Connect("127.0.0.1", 8080); while (true) { // 准备待监听的Socket集合(Select会修改集合,所以每次都要复制) var readSockets = new Socket[] { clientSocket }; Socket.Select(readSockets, null, null, -1); // 无限等待可读事件 // 处理服务器消息 if (readSockets.Length > 0 && readSockets[0] == clientSocket) { byte[] buffer = new byte[1024]; int bytesRead = clientSocket.Receive(buffer); if (bytesRead == 0) { Console.WriteLine("服务器已断开连接"); break; } Console.WriteLine($"服务器:{Encoding.UTF8.GetString(buffer, 0, bytesRead)}"); } // 处理控制台输入(检查是否有可用输入) if (Console.KeyAvailable) { string input = Console.ReadLine(); if (!string.IsNullOrWhiteSpace(input)) { byte[] sendBuffer = Encoding.UTF8.GetBytes(input); clientSocket.Send(sendBuffer); } } } } }
2. 使用System.Threading.Channels(现代异步编程模式)
利用Channel组件分离输入、接收、发送逻辑,代码结构更清晰,扩展性更强:
using System; using System.Net.Sockets; using System.Text; using System.Threading.Channels; using System.Threading.Tasks; class ChatClient { static async Task Main(string[] args) { // 连接服务器并获取网络流 using var clientSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); clientSocket.Connect("127.0.0.1", 8080); using var networkStream = new NetworkStream(clientSocket); // 创建无界通道传递用户输入 var inputChannel = Channel.CreateUnbounded<string>(); // 任务1:监听控制台输入,写入通道 var inputTask = Task.Run(async () => { while (true) { string input = Console.ReadLine(); if (!string.IsNullOrWhiteSpace(input)) { await inputChannel.Writer.WriteAsync(input); } } }); // 任务2:接收服务器消息并输出 var receiveTask = Task.Run(async () => { byte[] buffer = new byte[1024]; while (true) { int bytesRead = await networkStream.ReadAsync(buffer); if (bytesRead == 0) { Console.WriteLine("服务器已断开连接"); break; } Console.WriteLine($"服务器:{Encoding.UTF8.GetString(buffer, 0, bytesRead)}"); } }); // 任务3:读取通道中的用户输入,发送给服务器 var sendTask = Task.Run(async () => { await foreach (var input in inputChannel.Reader.ReadAllAsync()) { byte[] sendBuffer = Encoding.UTF8.GetBytes(input); await networkStream.WriteAsync(sendBuffer); } }); // 等待任一任务终止(比如服务器断开) await Task.WhenAny(receiveTask, inputTask, sendTask); } }
3. 使用TaskCompletionSource(底层异步控制)
通过TaskCompletionSource手动控制异步等待的触发,实现类似select的自定义多路复用:
using System; using System.Net.Sockets; using System.Text; using System.Threading.Tasks; class ChatClient { static async Task Main(string[] args) { using var clientSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); clientSocket.Connect("127.0.0.1", 8080); byte[] receiveBuffer = new byte[1024]; while (true) { // 创建两个TCS分别监听Socket接收和控制台输入 var socketReceiveTcs = new TaskCompletionSource<int>(); var consoleInputTcs = new TaskCompletionSource<string>(); // 启动Socket异步接收 _ = clientSocket.ReceiveAsync(receiveBuffer, SocketFlags.None).ContinueWith(t => { if (t.IsCompletedSuccessfully) socketReceiveTcs.SetResult(t.Result); else socketReceiveTcs.SetException(t.Exception ?? new Exception("接收失败")); }); // 启动控制台异步读取 _ = Task.Run(() => { string input = Console.ReadLine(); consoleInputTcs.SetResult(input); }); // 等待任一操作完成 var completedTask = await Task.WhenAny(socketReceiveTcs.Task, consoleInputTcs.Task); if (completedTask == socketReceiveTcs.Task) { int bytesRead = await socketReceiveTcs.Task; if (bytesRead == 0) { Console.WriteLine("服务器已断开连接"); break; } Console.WriteLine($"服务器:{Encoding.UTF8.GetString(receiveBuffer, 0, bytesRead)}"); } else if (completedTask == consoleInputTcs.Task) { string input = await consoleInputTcs.Task; if (!string.IsNullOrWhiteSpace(input)) { byte[] sendBuffer = Encoding.UTF8.GetBytes(input); await clientSocket.SendAsync(sendBuffer, SocketFlags.None); } } } } }
内容的提问来源于stack exchange,提问作者MwahMallah
相关产品推荐
相关产品推荐

