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