Pu Zhibing
2 天以前 f355ef485a56e613b71d0262c089b995d7ca10d2
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
package com.ruoyi.dataInterchange.util.mqtt;
 
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.channel.ExecutorChannel;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.messaging.MessageHandler;
 
import javax.annotation.Resource;
import java.util.concurrent.Executors;
 
/**
 * @author zhibing.pu
 * @Date 2025/5/23 17:26
 */
@Configuration
public class MqttConfiguration {
    
    @Resource
    private MqttConfigurationProperties mqttConfigurationProperties;
    
    
    
    /**
     * 连接参数
     *
     * @return
     */
    @Bean
    public MqttConnectOptions mqttConnectOptions() {
        MqttConnectOptions options = new MqttConnectOptions();
        options.setUserName(mqttConfigurationProperties.getUsername());
        options.setPassword(mqttConfigurationProperties.getPassword().toCharArray());
        options.setServerURIs(new String[]{mqttConfigurationProperties.getUrl()});
        options.setConnectionTimeout(mqttConfigurationProperties.getConnectionTimeout());
        options.setKeepAliveInterval(mqttConfigurationProperties.getKeepAlive());
        options.setCleanSession(true); // 设置为false以便断线重连后恢复会话
        options.setAutomaticReconnect(true);
        return options;
    }
    /**
     * 连接工厂
     *
     * @param options
     * @return
     */
    @Bean
    public MqttPahoClientFactory mqttClientFactory(MqttConnectOptions options) {
        DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
        factory.setConnectionOptions(options);
        return factory;
    }
    
    
 
    /**
     * 配置入站适配器
     *
     * @param mqttClientFactory
     * @return
     */
    @Bean
    public MqttPahoMessageDrivenChannelAdapter messageDrivenChannelAdapter(MqttPahoClientFactory mqttClientFactory) {
        MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter(mqttConfigurationProperties.getSubClientId(), mqttClientFactory);
        // adapter.addTopic("pub/300119110099"); 订阅主题,也可以放在初始化动态配置
        adapter.setOutputChannel(mqttInputChannel());
        return adapter;
    }
 
    /**
     * 支持多线程并发处理消息的输入通道
     *
     * @return
     */
    @Bean
    public ExecutorChannel mqttInputChannel() {
        return new ExecutorChannel(Executors.newFixedThreadPool(10)); // 线程池大小可以调整
    }
 
    /**
     * 配置消息处理器
     *
     * @return
     */
    @Bean
    @ServiceActivator(inputChannel = "mqttInputChannel") // 指定通道
    public MessageHandler messageHandler() {
        return new MqttReceiverMessageHandler();
    }
    
    
    
}