1、引入依赖

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-integration</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.integration</groupId>
    <artifactId>spring-integration-stream</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.integration</groupId>
    <artifactId>spring-integration-mqtt</artifactId>
</dependency>

2、配置项

mqtt:
  serverURIs: tcp://127.0.0.1:1883 #mqtt服务地址
  username: abc
  password: abc
  keepAliveInterval: 3        #心跳 单位 秒
  connectionTimeout: 5       #连接超时时间 单位 秒
  cleanSession: false         #是否清除session
  async: true                 #是否异步发送
  producer: #生产者配置
    defaultQos: 1              #消息质量 0:最多一次传输(可能丢包) 1:至少一次传输,(可能重包) 2:保证只有一次传输
    defaultRetained: false     #是否保留消息
    clientId: server_outbound  #生产者连接id
    defaultTopic: notification #默认订阅
  consumer: #消费者配置
    defaultTopic: notification #默认订阅
    clientId: server_inbound   #消费者连接id

3、MQTT配置类

@Configuration
@Slf4j
public class MqttConfig {
    @Value("${mqtt.serverURIs}")
    private String serverURIs;
    @Value("${mqtt.username}")
    private String username;
    @Value("${mqtt.password}")
    private String password;
    @Value("${mqtt.producer.clientId}")
    private String producerClientId;
    @Value("${mqtt.consumer.clientId}")
    private String consumerClientId;
    //心跳频率
    @Value("${mqtt.keepAliveInterval}")
    private int keepAliveInterval;
    //超时时间
    @Value("${mqtt.connectionTimeout}")
    private int connectionTimeout;
    //生产者默认订阅
    @Value("${mqtt.producer.defaultTopic}")
    private String producerDefaultTopic;
    //消费者默认订阅
    @Value("${mqtt.consumer.defaultTopic}")
    private String consumerDefaultTopic;
    //默认消息质量
    @Value("${mqtt.producer.defaultQos}")
    private int defaultProducerQos;
    //默认订阅
    @Value("${mqtt.producer.defaultRetained}")
    private boolean defaultRetained;

    //接收消息管道
    public static final String INBOUND_CHANNEL = "server_inbound";
    /**
     * Mqtt配置项
     * @return {@link org.eclipse.paho.client.mqttv3.MqttConnectOptions}
     */
    @Bean
    public MqttConnectOptions getMqttConnectOptions() {
        MqttConnectOptions options = new MqttConnectOptions();
        // 设置是否清空session,这里如果设置为false表示服务器会保留客户端的连接记录,
        // 这里设置为true表示每次连接到服务器都以新的身份连接
        options.setCleanSession(true);
        // 用户名
        options.setUserName(username);
        // 密码
        options.setPassword(password.toCharArray());
        // mqtt服务url
        options.setServerURIs(StringUtils.split(serverURIs, ","));
        // 设置超时时间 单位为秒
        options.setConnectionTimeout(connectionTimeout);
        // 设置会话心跳时间 单位为秒 服务器会每隔1.5*keepAliveInterval秒的时间向客户端发送心跳判断客户端是否在线,但这个方法并没有重连的机制
        options.setKeepAliveInterval(keepAliveInterval);
        // 设置“遗嘱”消息的话题,若客户端与服务器之间的连接意外中断,服务器将发布客户端的“遗嘱”消息。
        //options.setWill("willTopic", WILL_DATA, 2, false);
        return options;
    }

    /**
     * MQTT客户端
     * @return {@link org.springframework.integration.mqtt.core.MqttPahoClientFactory}
     */
    @Bean
    public MqttPahoClientFactory mqttClientFactory() {
        DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
        factory.setConnectionOptions(getMqttConnectOptions());
        return factory;
    }

    /**
     * mqtt消息管道(生产者)
     * @return {@link org.springframework.messaging.MessageChannel}
     */
    @Bean(name = IMqttSender.OUTBOUND_CHANNEL)
    public MessageChannel outboundChannel() {
        return new DirectChannel();
    }

    /**
     * mqtt消息处理器(生产者)
     * @return {@link org.springframework.messaging.MessageHandler}
     */
    @Bean
    @ServiceActivator(inputChannel = IMqttSender.OUTBOUND_CHANNEL)
    public MessageHandler getMqttProducer() {
        MqttPahoMessageHandler messageHandler = new MqttPahoMessageHandler(producerClientId, mqttClientFactory());
        messageHandler.setAsync(true);
        messageHandler.setDefaultTopic(producerDefaultTopic);
        //messageHandler.setDefaultRetained(defaultRetained);
        messageHandler.setDefaultQos(defaultProducerQos);
        return messageHandler;
    }


    /**
     * 消费者 如果不需要订阅消息,可不用配置下面
     */


    /**
     * MQTT信息通道(消费者)
     * @return {@link org.springframework.messaging.MessageChannel}
     */
    @Bean(name = INBOUND_CHANNEL)
    public MessageChannel mqttInboundChannel() {
        return new DirectChannel();
    }

    /**
     * MQTT消息订阅绑定(消费者)
     * @return {@link org.springframework.integration.core.MessageProducer}
     */
    @Bean
    public MessageProducer inbound() {
        // 可以同时消费(订阅)多个Topic
        MqttPahoMessageDrivenChannelAdapter adapter =
                new MqttPahoMessageDrivenChannelAdapter(
                        consumerClientId, mqttClientFactory(),
                        StringUtils.split(consumerDefaultTopic, ","));
        adapter.setCompletionTimeout(5000);
        adapter.setConverter(new DefaultPahoMessageConverter());
        adapter.setQos(1);
        // 设置订阅通道
        adapter.setOutputChannel(mqttInboundChannel());
        return adapter;
    }

    /**
     * MQTT消息处理器(消费者) 订阅的消息将会在这里打印
     *
     * @return {@link org.springframework.messaging.MessageHandler}
     */
    @Bean
    @ServiceActivator(inputChannel = INBOUND_CHANNEL)
    public MessageHandler handler() {
        return new MessageHandler() {
            @Override
            public void handleMessage(Message<?> message) throws MessagingException {
                log.info("--------------------接收到订阅消息--------------------");
                log.info("{}", message.getPayload());
            }
        };
    }
}

4、MQTT消息发送接口

@Component
@MessagingGateway(defaultRequestChannel = IMqttSender.OUTBOUND_CHANNEL)
public interface IMqttSender {
    //发送消息管道
    String OUTBOUND_CHANNEL = "server_outbound";
    /**
     * 发送信息到MQTT服务器
     * @param data 发送的文本
     */
    void sendToMqtt(String data);

    /**
     * 发送信息到MQTT服务器
     * @param topic 主题
     * @param payload 消息主体
     */
    void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic,
                    String payload);

    /**
     * 发送信息到MQTT服务器
     * @param topic 主题
     * @param qos 对消息处理的几种机制。<br> 0 表示的是订阅者没收到消息不会再次发送,消息会丢失。<br>
     * 1 表示的是会尝试重试,一直到接收到消息,但这种情况可能导致订阅者收到多次重复消息。<br>
     * 2 多了一次去重的动作,确保订阅者收到的消息有一次。
     * @param payload 消息主体
     */
    void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic,
                    @Header(MqttHeaders.QOS) int qos,
                    String payload);
}

注意事项:@MessagingGateway注解不会被@ComponentScan当做普通的组件扫描。放在Application同级目录及子目录时,会被默认扫描。如果需要将IMqttSender放在Application类的目录外,需要在Application类中配置@IntegrationComponentScan进行扫描。

5、发送消息到MQTT服务

@ActiveProfiles("test")
@SpringBootTest
public class MqttTest {
	@Resource
	private IMqttSender mqttSender;
	@Test
	void contextLoads(){
		mqttSender.sendToMqtt("测试发布消息");//默认topic 默认qos
		mqttSender.sendToMqtt("testTopic","测试发布消息");//自定义topic 默认qos
		mqttSender.sendToMqtt("testTopic",1,"测试发布消息");//自定义topic 自定义qos
	}
}
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐