From 8f3bf7050e65fdbe55eaad74fde307c57dab960e Mon Sep 17 00:00:00 2001
From: 云 <2163098428@qq.com>
Date: 星期五, 24 七月 2026 17:25:55 +0800
Subject: [PATCH] Merge remote-tracking branch 'origin/dev_business' into dev_business

---
 yudao-module-iot/yudao-module-iot-core/src/main/java/cn/iocoder/yudao/module/iot/core/messagebus/config/IotMessageBusAutoConfiguration.java |   51 ++++++++++++++++++++++++++++++++++++++++++++++++---
 1 files changed, 48 insertions(+), 3 deletions(-)

diff --git a/yudao-module-iot/yudao-module-iot-core/src/main/java/cn/iocoder/yudao/module/iot/core/messagebus/config/IotMessageBusAutoConfiguration.java b/yudao-module-iot/yudao-module-iot-core/src/main/java/cn/iocoder/yudao/module/iot/core/messagebus/config/IotMessageBusAutoConfiguration.java
index 67ae673..39f59f3 100644
--- a/yudao-module-iot/yudao-module-iot-core/src/main/java/cn/iocoder/yudao/module/iot/core/messagebus/config/IotMessageBusAutoConfiguration.java
+++ b/yudao-module-iot/yudao-module-iot-core/src/main/java/cn/iocoder/yudao/module/iot/core/messagebus/config/IotMessageBusAutoConfiguration.java
@@ -6,7 +6,9 @@
 import cn.iocoder.yudao.framework.mq.redis.core.stream.AbstractRedisStreamMessage;
 import cn.iocoder.yudao.framework.mq.redis.core.stream.AbstractRedisStreamMessageListener;
 import cn.iocoder.yudao.module.iot.core.messagebus.core.IotMessageBus;
+import cn.iocoder.yudao.module.iot.core.messagebus.core.kafka.IotKafkaMessageBus;
 import cn.iocoder.yudao.module.iot.core.messagebus.core.local.IotLocalMessageBus;
+import cn.iocoder.yudao.module.iot.core.messagebus.core.rabbitmq.IotRabbitMQMessageBus;
 import cn.iocoder.yudao.module.iot.core.messagebus.core.redis.IotRedisMessageBus;
 import cn.iocoder.yudao.module.iot.core.messagebus.core.rocketmq.IotRocketMQMessageBus;
 import cn.iocoder.yudao.module.iot.core.mq.producer.IotDeviceMessageProducer;
@@ -14,15 +16,20 @@
 import org.apache.rocketmq.spring.autoconfigure.RocketMQProperties;
 import org.apache.rocketmq.spring.core.RocketMQTemplate;
 import org.redisson.api.RedissonClient;
+import org.springframework.amqp.rabbit.core.RabbitAdmin;
+import org.springframework.amqp.rabbit.core.RabbitTemplate;
 import org.springframework.boot.autoconfigure.AutoConfiguration;
 import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
 import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
 import org.springframework.boot.context.properties.EnableConfigurationProperties;
+import org.springframework.boot.kafka.autoconfigure.KafkaProperties;
 import org.springframework.context.ApplicationContext;
 import org.springframework.context.annotation.Bean;
 import org.springframework.context.annotation.Configuration;
 import org.springframework.data.redis.core.RedisTemplate;
 import org.springframework.data.redis.core.StringRedisTemplate;
+import org.springframework.kafka.core.KafkaTemplate;
 
 import java.util.List;
 
@@ -73,6 +80,21 @@
 
     }
 
+    // ==================== Kafka 瀹炵幇 ====================
+
+    @Configuration
+    @ConditionalOnProperty(prefix = "yudao.iot.message-bus", name = "type", havingValue = "kafka")
+    @ConditionalOnClass(KafkaTemplate.class)
+    public static class IotKafkaMessageBusConfiguration {
+
+        @Bean
+        public IotKafkaMessageBus iotKafkaMessageBus(KafkaProperties kafkaProperties) {
+            log.info("[iotKafkaMessageBus][鍒涘缓 IoT Kafka 娑堟伅鎬荤嚎]");
+            return new IotKafkaMessageBus(kafkaProperties);
+        }
+
+    }
+
     // ==================== Redis 瀹炵幇 ====================
 
     /**
@@ -99,7 +121,8 @@
                                                                             RedisMQTemplate redisTemplate,
                                                                             RedissonClient redissonClient) {
             List<AbstractRedisStreamMessageListener<?>> listeners = getListeners(messageBus);
-            return new RedisPendingMessageResendJob(listeners, redisTemplate, redissonClient);
+            return new RedisPendingMessageResendJob(listeners, redisTemplate, redissonClient,
+                    RedisPendingMessageResendJob.IOT_RESEND_LOCK_KEY);
         }
 
         /**
@@ -110,7 +133,8 @@
                                                                             RedisMQTemplate redisTemplate,
                                                                             RedissonClient redissonClient) {
             List<AbstractRedisStreamMessageListener<?>> listeners = getListeners(messageBus);
-            return new RedisStreamMessageCleanupJob(listeners, redisTemplate, redissonClient);
+            return new RedisStreamMessageCleanupJob(listeners, redisTemplate, redissonClient,
+                    RedisStreamMessageCleanupJob.IOT_CLEANUP_LOCK_KEY);
         }
 
         private List<AbstractRedisStreamMessageListener<?>> getListeners(IotRedisMessageBus messageBus) {
@@ -126,4 +150,25 @@
 
     }
 
-}
\ No newline at end of file
+    // ==================== RabbitMQ 瀹炵幇 ====================
+
+    @Configuration
+    @ConditionalOnProperty(prefix = "yudao.iot.message-bus", name = "type", havingValue = "rabbitmq")
+    @ConditionalOnClass(RabbitTemplate.class)
+    public static class IotRabbitMQMessageBusConfiguration {
+
+        @Bean
+        @ConditionalOnMissingBean
+        public RabbitAdmin rabbitAdmin(RabbitTemplate rabbitTemplate) {
+            return new RabbitAdmin(rabbitTemplate);
+        }
+
+        @Bean
+        public IotRabbitMQMessageBus iotRabbitMQMessageBus(RabbitTemplate rabbitTemplate, RabbitAdmin rabbitAdmin) {
+            log.info("[iotRabbitMQMessageBus][鍒涘缓 IoT RabbitMQ 娑堟伅鎬荤嚎]");
+            return new IotRabbitMQMessageBus(rabbitTemplate, rabbitAdmin);
+        }
+
+    }
+
+}

--
Gitblit v1.9.3