如何在C#中连接两个控制台应用并实现循环队列数据交互?
当然有可行的方案!首先得明确一个关键点:两个独立的控制台应用是完全隔离的进程,它们的内存空间不共享,所以你没办法让B直接访问A里的StaticQueue静态队列。必须通过**跨进程通信(IPC)**机制来让A把队列里的数据传递给B,然后B再处理自己的LocalQueue逻辑。
可行方案概述
核心思路是用IPC实现A和B之间的数据传递,同时处理两边队列的业务逻辑。我推荐用.NET原生支持的命名管道(Named Pipes),它轻量、适合本地进程间通信,而且能很好地处理消息边界,避免数据粘包问题。
具体实现步骤
首先我们需要实现一个线程安全的循环队列(因为C#内置的Queue<T>不是循环队列,而且多线程访问需要同步),然后分别编写A和B的代码。
1. 线程安全的循环队列实现
先写一个通用的循环队列类,确保多线程访问时不会出现数据错乱:
public class CircularQueue<T> { private readonly T[] _items; private int _head; private int _tail; private int _count; private readonly object _lockObj = new object(); public int Capacity => _items.Length; public int Count { get { lock (_lockObj) return _count; } } public bool IsEmpty { get { lock (_lockObj) return _count == 0; } } public bool IsFull { get { lock (_lockObj) return _count == _items.Length; } } public CircularQueue(int capacity) { _items = new T[capacity]; _head = _tail = _count = 0; } public bool Enqueue(T item) { lock (_lockObj) { if (_count == _items.Length) return false; _items[_tail] = item; _tail = (_tail + 1) % _items.Length; _count++; return true; } } public bool TryDequeue(out T item) { lock (_lockObj) { if (_count == 0) { item = default; return false; } item = _items[_head]; _head = (_head + 1) % _items.Length; _count--; return true; } } }
2. 控制台应用A的代码
A需要做两件事:持续填充自己的StaticQueue,同时作为命名管道的服务器,把队列里的数发送给B:
using System.IO.Pipes; using System.Text.Json; namespace ConsoleAppA { class Program { private static readonly CircularQueue<int> StaticQueue = new CircularQueue<int>(100); // 设定A的队列容量 private static readonly CancellationTokenSource Cts = new CancellationTokenSource(); static async Task Main(string[] args) { // 后台任务:持续填充队列 _ = FillQueueAsync(Cts.Token); // 启动命名管道服务器,等待B连接并发送数据 await RunPipeServerAsync(Cts.Token); Console.WriteLine("按任意键退出..."); Console.ReadKey(); Cts.Cancel(); } private static async Task FillQueueAsync(CancellationToken token) { var random = new Random(); while (!token.IsCancellationRequested) { if (!StaticQueue.IsFull) { int num = random.Next(1, 100); StaticQueue.Enqueue(num); Console.WriteLine($"A: 已加入队列 -> {num},当前队列元素数: {StaticQueue.Count}/{StaticQueue.Capacity}"); } await Task.Delay(500, token); // 模拟填充间隔 } } private static async Task RunPipeServerAsync(CancellationToken token) { while (!token.IsCancellationRequested) { try { // 创建命名管道服务器,管道名为"QueuePipe" using var pipeServer = new NamedPipeServerStream("QueuePipe", PipeDirection.Out, 1, PipeTransmissionMode.Message); Console.WriteLine("A: 等待B连接..."); await pipeServer.WaitForConnectionAsync(token); Console.WriteLine("A: B已连接"); // 逐个发送队列中的元素,直到队列空、B断开或任务取消 while (!token.IsCancellationRequested && pipeServer.IsConnected && !StaticQueue.IsEmpty) { if (StaticQueue.TryDequeue(out int num)) { // 序列化数据为JSON(也可以用二进制,更高效) var json = JsonSerializer.Serialize(num); var buffer = System.Text.Encoding.UTF8.GetBytes(json); await pipeServer.WriteAsync(buffer, 0, buffer.Length, token); await pipeServer.FlushAsync(token); Console.WriteLine($"A: 已发送数据 -> {num}"); await Task.Delay(100, token); // 模拟发送间隔 } } pipeServer.Disconnect(); Console.WriteLine("A: B已断开连接"); } catch (OperationCanceledException) { break; } catch (Exception ex) { Console.WriteLine($"A: 管道异常 -> {ex.Message}"); await Task.Delay(1000, token); } } } } }
3. 控制台应用B的代码
B需要连接到A的命名管道,接收数据,然后按照需求将数据加入自己的LocalQueue(仅当LocalQueue非空且未达到容量时加入):
using System.IO.Pipes; using System.Text.Json; namespace ConsoleAppB { class Program { private static readonly CircularQueue<int> LocalQueue = new CircularQueue<int>(20); // 设定B的队列容量 private static readonly CancellationTokenSource Cts = new CancellationTokenSource(); static async Task Main(string[] args) { // 后台任务:从A拉取数据并处理 _ = PullDataFromAAsync(Cts.Token); Console.WriteLine("按任意键退出..."); Console.ReadKey(); Cts.Cancel(); } private static async Task PullDataFromAAsync(CancellationToken token) { while (!token.IsCancellationRequested) { try { // 连接到A的命名管道 using var pipeClient = new NamedPipeClientStream(".", "QueuePipe", PipeDirection.In); Console.WriteLine("B: 正在连接A..."); await pipeClient.ConnectAsync(5000, token); Console.WriteLine("B: 已连接到A"); var buffer = new byte[1024]; while (!token.IsCancellationRequested && pipeClient.IsConnected) { int bytesRead = await pipeClient.ReadAsync(buffer, 0, buffer.Length, token); if (bytesRead == 0) break; // 连接断开 // 反序列化数据 var json = System.Text.Encoding.UTF8.GetString(buffer, 0, bytesRead); if (JsonSerializer.Deserialize<int>(json) is int num) { // 严格按照需求:检查LocalQueue是否非空,若非空且未达容量则加入 if (!LocalQueue.IsEmpty && !LocalQueue.IsFull) { bool added = LocalQueue.Enqueue(num); Console.WriteLine($"B: 已加入LocalQueue -> {num},当前元素数: {LocalQueue.Count}/{LocalQueue.Capacity}"); } else { if (LocalQueue.IsEmpty) { Console.WriteLine($"B: LocalQueue为空,跳过加入 -> {num}"); } else if (LocalQueue.IsFull) { Console.WriteLine($"B: LocalQueue已满,停止接收数据"); break; // 已满,断开连接 } } } } Console.WriteLine("B: 与A的连接已断开"); } catch (OperationCanceledException) { break; } catch (TimeoutException) { Console.WriteLine("B: 连接超时,1秒后重试..."); await Task.Delay(1000, token); } catch (Exception ex) { Console.WriteLine($"B: 管道异常 -> {ex.Message}"); await Task.Delay(1000, token); } } } } }
其他可选方案
如果你的需求更复杂,还可以考虑这些IPC方式:
- 内存映射文件:适合共享大量数据,需要自己用
Mutex处理同步逻辑。 - MSMQ(微软消息队列):更重量级,支持离线消息和分布式场景,但配置稍繁琐。
- gRPC:适合跨机器通信,或者需要更复杂的接口定义,但本地进程用它有点“杀鸡用牛刀”。
注意事项
- 线程安全:A的队列被两个后台任务同时操作,必须加锁(我们的
CircularQueue已经处理了这一点)。 - 异常处理:要处理连接超时、断开、取消等异常情况,保证程序的稳定性。
- 消息边界:命名管道的
PipeTransmissionMode.Message模式确保我们能完整接收每一条数据,避免粘包问题。
内容的提问来源于stack exchange,提问作者the_coder_guy
相关产品推荐
相关产品推荐

