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

求助:基于传感器列表动态管理RabbitMQ生产者消费者的C#实现

问题需求

我需要在运行时基于传感器列表动态创建和销毁RabbitMQ队列与消费者,目前编写的代码仅能部分正常工作,恳请有经验的开发者协助完成多传感器及监控代理的相关功能开发。

现有代码

消费者类

public class SensorQueueConsumer
{
    public static void Consume(IModel channel)
    {
        Tasks.AddData();
        foreach (string sensor in Sensor.SensorList)
        {
            channel.QueueDeclare($"{sensor}Queue",
                durable: true,
                exclusive: false,
                autoDelete: false,
                arguments: null);
            var consumer = new EventingBasicConsumer(channel);
            consumer.Received += (sender, e) =>
            {
                var body = e.Body.ToArray();
                var message = Encoding.UTF8.GetString(body);
                Console.WriteLine(message);
            };

            channel.BasicConsume($"{sensor}Queue", true, consumer);
            Console.WriteLine("consumer started");
            Console.ReadLine();
        }
    }
}

生产者类

public static class SensorQueueProducer
{
    public static void Publish(IModel channel)
    {
        Tasks.AddData();
        foreach (string sensor in Sensor.SensorList)
        {
            channel.QueueDeclare($"{sensor}Queue",
                durable: true,
                exclusive: false,
                autoDelete: false,
                arguments: null);
            var Count = 0;
            while (true)
            {
                var message = new { Name = "PlantEvent", Sensor = sensor, Saverity = "high", Message = $"Hello! Water me!! Count:{Count}" };
                var body = Encoding.UTF8.GetBytes(JsonConvert.SerializeObject(message));

                channel.BasicPublish("", $"{sensor}Queue", null, body);
                Count++;
                Thread.Sleep(1000);
            }
        }
    }
}

传感器列表类

public class Sensor
{
    public static List<String> SensorList { get; set; } = new List<String>();
}

public class Tasks
{
    public static void AddData()
    {
        string[] sensors = { "ec", "soil", "co2" };
        Sensor.SensorList.AddRange(sensors);
    }

    public static void ClearData()
    {
        Sensor.SensorList.Clear();
    }
}

现有代码的核心问题

  • 消费者阻塞:Console.ReadLine()会终止循环,导致仅能初始化第一个传感器的消费者,后续传感器队列无法创建和启动消费。
  • 生产者阻塞:while(true)死循环会卡在第一个传感器的发布逻辑,无法处理后续传感器的消息推送。
  • 无动态销毁逻辑:没有实现队列删除、消费者取消的功能,无法响应传感器列表的动态变更。
  • 静态类设计缺陷:静态列表和静态方法无法有效管理消费者、生产者的状态,比如无法跟踪已启动的消费者标签,后续无法取消消费。

改进后的实现方案

1. 重构传感器列表(线程安全)

public class SensorManager
{
    private readonly List<string> _sensorList = new List<string>();
    private readonly object _lockObj = new object();

    public IReadOnlyList<string> SensorList
    {
        get
        {
            lock (_lockObj)
            {
                return _sensorList.ToList();
            }
        }
    }

    public void AddSensors(IEnumerable<string> sensors)
    {
        lock (_lockObj)
        {
            _sensorList.AddRange(sensors.Except(_sensorList));
        }
    }

    public void RemoveSensor(string sensor)
    {
        lock (_lockObj)
        {
            _sensorList.Remove(sensor);
        }
    }

    public void ClearSensors()
    {
        lock (_lockObj)
        {
            _sensorList.Clear();
        }
    }
}

2. 动态消费者管理器

public class SensorConsumerManager : IDisposable
{
    private readonly IModel _channel;
    private readonly SensorManager _sensorManager;
    private readonly Dictionary<string, string> _consumerTags = new Dictionary<string, string>();
    private readonly object _lockObj = new object();

    public SensorConsumerManager(IModel channel, SensorManager sensorManager)
    {
        _channel = channel;
        _sensorManager = sensorManager;
    }

    // 启动所有传感器的消费者
    public void StartAllConsumers()
    {
        foreach (var sensor in _sensorManager.SensorList)
        {
            StartConsumerForSensor(sensor);
        }
    }

    // 为单个传感器启动消费者
    public void StartConsumerForSensor(string sensor)
    {
        lock (_lockObj)
        {
            if (_consumerTags.ContainsKey(sensor))
                return;

            var queueName = $"{sensor}Queue";
            _channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false, arguments: null);

            var consumer = new EventingBasicConsumer(_channel);
            consumer.Received += (sender, e) =>
            {
                var body = e.Body.ToArray();
                var message = Encoding.UTF8.GetString(body);
                Console.WriteLine($"[{sensor}] 收到消息: {message}");
            };

            var consumerTag = _channel.BasicConsume(queueName, autoAck: true, consumer: consumer);
            _consumerTags.Add(sensor, consumerTag);
            Console.WriteLine($"[{sensor}] 消费者已启动");
        }
    }

    // 停止单个传感器的消费者并删除队列(可选)
    public void StopConsumerForSensor(string sensor, bool deleteQueue = false)
    {
        lock (_lockObj)
        {
            if (!_consumerTags.TryGetValue(sensor, out var tag))
                return;

            _channel.BasicCancel(tag);
            _consumerTags.Remove(sensor);
            Console.WriteLine($"[{sensor}] 消费者已停止");

            if (deleteQueue)
            {
                var queueName = $"{sensor}Queue";
                _channel.QueueDelete(queueName);
                Console.WriteLine($"[{sensor}] 队列已删除");
            }
        }
    }

    // 停止所有消费者
    public void StopAllConsumers(bool deleteQueues = false)
    {
        foreach (var sensor in _consumerTags.Keys.ToList())
        {
            StopConsumerForSensor(sensor, deleteQueues);
        }
    }

    public void Dispose()
    {
        StopAllConsumers();
        _channel?.Close();
        _channel?.Dispose();
    }
}

3. 动态生产者管理器

public class SensorProducerManager : IDisposable
{
    private readonly IModel _channel;
    private readonly SensorManager _sensorManager;
    private readonly Dictionary<string, CancellationTokenSource> _producerTokens = new Dictionary<string, CancellationTokenSource>();
    private readonly object _lockObj = new object();

    public SensorProducerManager(IModel channel, SensorManager sensorManager)
    {
        _channel = channel;
        _sensorManager = sensorManager;
    }

    // 启动所有传感器的生产者
    public void StartAllProducers()
    {
        foreach (var sensor in _sensorManager.SensorList)
        {
            StartProducerForSensor(sensor);
        }
    }

    // 为单个传感器启动生产者
    public void StartProducerForSensor(string sensor)
    {
        lock (_lockObj)
        {
            if (_producerTokens.ContainsKey(sensor))
                return;

            var cts = new CancellationTokenSource();
            _producerTokens.Add(sensor, cts);

            Task.Run(async () =>
            {
                var queueName = $"{sensor}Queue";
                _channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false, arguments: null);
                var count = 0;

                try
                {
                    while (!cts.Token.IsCancellationRequested)
                    {
                        var message = new
                        {
                            Name = "PlantEvent",
                            Sensor = sensor,
                            Severity = "high",
                            Message = $"Hello! Water me!! Count:{count}"
                        };
                        var body = Encoding.UTF8.GetBytes(JsonConvert.SerializeObject(message));

                        _channel.BasicPublish("", queueName, null, body);
                        Console.WriteLine($"[{sensor}] 推送消息: {message.Message}");
                        count++;
                        await Task.Delay(1000, cts.Token);
                    }
                }
                catch (TaskCanceledException)
                {
                    Console.WriteLine($"[{sensor}] 生产者已停止");
                }
            }, cts.Token);
        }
    }

    // 停止单个传感器的生产者
    public void StopProducerForSensor(string sensor)
    {
        lock (_lockObj)
        {
            if (_producerTokens.TryGetValue(sensor, out var cts))
            {
                cts.Cancel();
                _producerTokens.Remove(sensor);
            }
        }
    }

    // 停止所有生产者
    public void StopAllProducers()
    {
        foreach (var cts in _producerTokens.Values)
        {
            cts.Cancel();
        }
        _producerTokens.Clear();
    }

    public void Dispose()
    {
        StopAllProducers();
        _channel?.Close();
        _channel?.Dispose();
    }
}

4. 使用示例

class Program
{
    static async Task Main(string[] args)
    {
        var factory = new ConnectionFactory() { HostName = "localhost" };
        using var connection = factory.CreateConnection();

        var sensorManager = new SensorManager();
        sensorManager.AddSensors(new[] { "ec", "soil", "co2" });

        using var consumerChannel = connection.CreateModel();
        var consumerManager = new SensorConsumerManager(consumerChannel, sensorManager);
        consumerManager.StartAllConsumers();

        using var producerChannel = connection.CreateModel();
        var producerManager = new SensorProducerManager(producerChannel, sensorManager);
        producerManager.StartAllProducers();

        // 模拟动态移除传感器
        await Task.Delay(5000);
        Console.WriteLine("\n开始移除soil传感器...");
        sensorManager.RemoveSensor("soil");
        consumerManager.StopConsumerForSensor("soil", deleteQueue: true);
        producerManager.StopProducerForSensor("soil");

        // 模拟新增传感器
        await Task.Delay(3000);
        Console.WriteLine("\n新增temp传感器...");
        sensorManager.AddSensors(new[] { "temp" });
        consumerManager.StartConsumerForSensor("temp");
        producerManager.StartProducerForSensor("temp");

        Console.WriteLine("\n按任意键退出...");
        Console.ReadKey();
    }
}

关键改进点

  • 线程安全:所有共享状态操作加锁,避免多线程冲突。
  • 非阻塞设计:生产者使用Task和取消令牌实现异步推送,消费者启动后不阻塞主线程。
  • 动态管理:支持单独启停单个传感器的消费/生产,可选择是否删除队列。
  • 资源自动清理:实现IDisposable接口,确保程序退出时释放RabbitMQ资源。
  • 错误修正:修复了原代码中Saverity的拼写错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 18:31:16