You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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:适合跨机器通信,或者需要更复杂的接口定义,但本地进程用它有点“杀鸡用牛刀”。
注意事项
  1. 线程安全:A的队列被两个后台任务同时操作,必须加锁(我们的CircularQueue已经处理了这一点)。
  2. 异常处理:要处理连接超时、断开、取消等异常情况,保证程序的稳定性。
  3. 消息边界:命名管道的PipeTransmissionMode.Message模式确保我们能完整接收每一条数据,避免粘包问题。

内容的提问来源于stack exchange,提问作者the_coder_guy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.30 06:32:37