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
99
package cn.iocoder.yudao.module.iot.core.messagebus.core.redis;
 
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.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.connection.stream.*;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.data.redis.stream.StreamMessageListenerContainer;
 
import java.lang.reflect.Type;
import java.util.ArrayList;
import java.util.List;
 
import static cn.iocoder.yudao.framework.mq.redis.config.YudaoRedisMQConsumerAutoConfiguration.buildConsumerName;
import static cn.iocoder.yudao.framework.mq.redis.config.YudaoRedisMQConsumerAutoConfiguration.checkRedisVersion;
 
/**
 * Redis 的 {@link IotMessageBus} 实现类
 *
 * @author 芋道源码
 */
@Slf4j
public class IotRedisMessageBus implements IotMessageBus {
 
    private final RedisTemplate<String, ?> redisTemplate;
 
    private final StreamMessageListenerContainer<String, ObjectRecord<String, String>> redisStreamMessageListenerContainer;
 
    @Getter
    private final List<IotMessageSubscriber<?>> subscribers = new ArrayList<>();
 
    public IotRedisMessageBus(RedisTemplate<String, ?> redisTemplate) {
        this.redisTemplate = redisTemplate;
        checkRedisVersion(redisTemplate);
        // 创建 options 配置
        StreamMessageListenerContainer.StreamMessageListenerContainerOptions<String, ObjectRecord<String, String>> containerOptions =
                StreamMessageListenerContainer.StreamMessageListenerContainerOptions.builder()
                        .batchSize(10) // 一次性最多拉取多少条消息
                        .targetType(String.class) // 目标类型。统一使用 String,通过自己封装的 AbstractStreamMessageListener 去反序列化
                        .build();
        // 创建 container 对象
        this.redisStreamMessageListenerContainer =
                StreamMessageListenerContainer.create(redisTemplate.getRequiredConnectionFactory(), containerOptions);
    }
 
    @PostConstruct
    public void init() {
        this.redisStreamMessageListenerContainer.start();
    }
 
    @PreDestroy
    public void destroy() {
        this.redisStreamMessageListenerContainer.stop();
    }
 
    @Override
    public void post(String topic, Object message) {
        redisTemplate.opsForStream().add(StreamRecords.newRecord()
                .ofObject(JsonUtils.toJsonString(message)) // 设置内容
                .withStreamKey(topic)); // 设置 stream key
    }
 
    @Override
    public void register(IotMessageSubscriber<?> subscriber) {
        Type type = TypeUtil.getTypeArgument(subscriber.getClass(), 0);
        if (type == null) {
            throw new IllegalStateException(String.format("类型(%s) 需要设置消息类型", getClass().getName()));
        }
 
        // 创建 listener 对应的消费者分组
        try {
            redisTemplate.opsForStream().createGroup(subscriber.getTopic(), subscriber.getGroup());
        } catch (Exception ignore) {
        }
        // 创建 Consumer 对象
        String consumerName = buildConsumerName();
        Consumer consumer = Consumer.from(subscriber.getGroup(), consumerName);
        // 设置 Consumer 消费进度,以最小消费进度为准
        StreamOffset<String> streamOffset = StreamOffset.create(subscriber.getTopic(), ReadOffset.lastConsumed());
        // 设置 Consumer 监听
        StreamMessageListenerContainer.StreamReadRequestBuilder<String> builder = StreamMessageListenerContainer.StreamReadRequest
                .builder(streamOffset).consumer(consumer)
                .autoAcknowledge(false) // 不自动 ack
                .cancelOnError(throwable -> false); // 默认配置,发生异常就取消消费,显然不符合预期;因此,我们设置为 false
        redisStreamMessageListenerContainer.register(builder.build(), message -> {
            // 消费消息
            subscriber.onMessage(JsonUtils.parseObject(message.getValue(), type));
            // ack 消息消费完成
            redisTemplate.opsForStream().acknowledge(subscriber.getGroup(), message);
        });
        this.subscribers.add(subscriber);
    }
 
}