MqttTopicManager.java
2.47 KB
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
package com.diligrp.mqtt.integration.manager;
import com.diligrp.mqtt.core.model.SubscribeModel;
import com.diligrp.mqtt.core.service.MqttTopicService;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.stereotype.Service;
import java.util.Arrays;
import java.util.HashSet;
import java.util.Optional;
/**
* MQTT动态Topic管理服务
*/
@Slf4j
@Service
public class MqttTopicManager implements MqttTopicService {
@Resource
private MqttPahoMessageDrivenChannelAdapter mqttMessageDrivenChannelAdapter;
/**
* 订阅主题
*
* @param subscribeModel 订阅模式
*/
@Override
public void subscribeTopic(SubscribeModel subscribeModel) {
subscribeTopic(subscribeModel.getTopic(), subscribeModel.getQos());
}
/**
* 动态订阅Topic
*
* @param topic 要订阅的Topic
* @param qos QoS级别
*/
@Override
public void subscribeTopic(String topic, int qos) {
Optional.ofNullable(topic).filter(t -> !isTopicSubscribed(t)).ifPresentOrElse(t -> {
mqttMessageDrivenChannelAdapter.addTopic(t, qos);
log.info("成功订阅Topic: {}, QoS: {}", t, qos);
}, () -> log.info("无效订阅Topic: {}, QoS: {}", topic, qos));
}
/**
* 取消订阅Topic
*
* @param topic 要取消的Topic
*/
@Override
public void unsubscribeTopic(String topic) {
Optional.ofNullable(topic).filter(this::isTopicSubscribed).ifPresentOrElse(t -> {
mqttMessageDrivenChannelAdapter.removeTopic(t);
log.info("成功取消订阅Topic: {}", topic);
}, () -> log.info("无效取消订阅Topic: {}", topic));
}
/**
* 检查Topic是否已订阅
*
* @param topic 要检查的Topic
* @return 是否订阅
*/
public boolean isTopicSubscribed(String topic) {
return Optional.ofNullable(mqttMessageDrivenChannelAdapter.getTopic()).filter(t -> t.length > 0).map(ts -> new HashSet<>(Arrays.asList(ts))).orElseGet(HashSet::new).contains(topic);
}
/**
* 获取所有已订阅的Topic
*
* @return 订阅的Topic数组
*/
@Override
public HashSet<String> getSubscribedTopics() {
return Optional.ofNullable(mqttMessageDrivenChannelAdapter.getTopic()).filter(t -> t.length > 0).map(ts -> new HashSet<>(Arrays.asList(ts))).orElseGet(HashSet::new);
}
}