求助:基于传感器列表动态管理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
相关产品推荐
相关产品推荐

