2026-06-26 20b96473f2520590a0dca6b775b81e3ea06a77a0
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
package cn.iocoder.yudao.module.iot.core.messagebus.core.rocketmq;
 
import cn.hutool.core.util.TypeUtil;
import cn.iocoder.yudao.framework.common.util.json.JsonUtils;
import cn.iocoder.yudao.module.iot.core.messagebus.core.IotMessageBus;
import cn.iocoder.yudao.module.iot.core.messagebus.core.IotMessageSubscriber;
import jakarta.annotation.PreDestroy;
import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.spring.autoconfigure.RocketMQProperties;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
 
import java.lang.reflect.Type;
import java.util.ArrayList;
import java.util.List;
 
/**
 * 基于 RocketMQ 的 {@link IotMessageBus} 实现类
 *
 * @author 芋道源码
 */
@RequiredArgsConstructor
@Slf4j
public class IotRocketMQMessageBus implements IotMessageBus {
 
    private final RocketMQProperties rocketMQProperties;
 
    private final RocketMQTemplate rocketMQTemplate;
 
    /**
     * 主题对应的消费者映射
     */
    private final List<DefaultMQPushConsumer> topicConsumers = new ArrayList<>();
 
    /**
     * 销毁时关闭所有消费者
     */
    @PreDestroy
    public void destroy() {
        for (DefaultMQPushConsumer consumer : topicConsumers) {
            try {
                consumer.shutdown();
                log.info("[destroy][关闭 group({}) 的消费者成功]", consumer.getConsumerGroup());
            } catch (Exception e) {
                log.error("[destroy]关闭 group({}) 的消费者异常]", consumer.getConsumerGroup(), e);
            }
        }
    }
 
    @Override
    public void post(String topic, Object message) {
        // TODO @芋艿:需要 orderly!
        SendResult result = rocketMQTemplate.syncSend(topic, JsonUtils.toJsonString(message));
        log.info("[post][topic({}) 发送消息({}) result({})]", topic, message, result);
    }
 
    @Override
    @SneakyThrows
    public void register(IotMessageSubscriber<?> subscriber) {
        Type type = TypeUtil.getTypeArgument(subscriber.getClass(), 0);
        if (type == null) {
            throw new IllegalStateException(String.format("类型(%s) 需要设置消息类型", getClass().getName()));
        }
 
        // 1.1 创建 DefaultMQPushConsumer
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer();
        consumer.setNamesrvAddr(rocketMQProperties.getNameServer());
        consumer.setConsumerGroup(subscriber.getGroup());
        // 1.2 订阅主题
        consumer.subscribe(subscriber.getTopic(), "*");
        // 1.3 设置消息监听器
        consumer.setMessageListener((MessageListenerConcurrently) (messages, context) -> {
            for (MessageExt messageExt : messages) {
                try {
                    byte[] body = messageExt.getBody();
                    subscriber.onMessage(JsonUtils.parseObject(body, type));
                } catch (Exception ex) {
                    log.error("[onMessage][topic({}/{}) message({}) 消费者({}) 处理异常]",
                            subscriber.getTopic(), subscriber.getGroup(), messageExt, subscriber.getClass().getName(), ex);
                    return ConsumeConcurrentlyStatus.RECONSUME_LATER;
                }
            }
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        });
        // 1.4 启动消费者
        consumer.start();
 
        // 2. 保存消费者引用
        topicConsumers.add(consumer);
    }
 
}