最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot 集成MQTT實現(xiàn)消息訂閱的詳細代碼

 更新時間:2024年11月28日 14:35:49   作者:不甘平凡--liang  
本文介紹了如何在SpringBoot中集成MQTT并實現(xiàn)消息訂閱,主要步驟包括添加依賴、配置文件設置、啟動類注解、MQTT配置類、消息處理器配置、主題緩存、動態(tài)數(shù)據(jù)庫主題配置以及消息處理服務,感興趣的朋友跟隨小編一起看看吧

1、引入依賴

  <!--MQTT start-->
 <dependency>
      <groupId>org.springframework.boot</groupId>
      <artifactId>spring-boot-starter-integration</artifactId>
  </dependency>
  <dependency>
     <groupId>org.springframework.integration</groupId>
     <artifactId>spring-integration-mqtt</artifactId>
     <version>5.4.4</version>
 </dependency>
 <!--MQTT end-->
 <dependency>
     <groupId>org.springframework.boot</groupId>
      <artifactId>spring-boot-configuration-processor</artifactId>
      <optional>true</optional>
  </dependency>

2、增加yml配置

  spring:
    mqtt:
      username: test
      password: test
      url: tcp://127.0.0.1:8080
      subClientId: singo_sub_client_id_888 #訂閱 客戶端id
      pubClientId: singo_pub_client_id_888 #發(fā)布 客戶端id
      connectionTimeout: 30
      keepAlive: 60

3、資源配置類

@Data
@ConfigurationProperties(prefix = "spring.mqtt")
public class MqttConfigurationProperties {
    private String username;
    private String password;
    private String url;
    private String subClientId;
    private String pubClientId;
    private int connectionTimeout;
    private int keepAlive;
}

注意啟動類需要增加注解

@EnableConfigurationProperties(MqttConfigurationProperties.class)

4、MQTT配置類

@Configuration
public class MqttConfig {
    @Autowired
    private MqttConfigurationProperties mqttConfigurationProperties;
    /**
     * 連接參數(shù)
     *
     * @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;
    }
    /**
     * 消息輸入通道
     * 每次只有一個消息處理器可以消費消息。
     * 當前消息的處理完成之前,新消息需要排隊等待,無法并行處理。
     * 默認是:單線程、順序執(zhí)行的
     * @return
     */
    // @Bean
    // public DirectChannel mqttInputChannel() {
    //     return new DirectChannel();
    // }
    /**
     * 支持多線程并發(fā)處理消息的輸入通道
     *
     * @return
     */
    @Bean
    public ExecutorChannel mqttInputChannel() {
        return new ExecutorChannel(Executors.newFixedThreadPool(10)); // 線程池大小可以調(diào)整
    }
    /**
     * 配置入站適配器
     *
     * @param mqttClientFactory
     * @return
     */
    @Bean
    public MqttPahoMessageDrivenChannelAdapter messageDrivenChannelAdapter(MqttPahoClientFactory mqttClientFactory) {
        MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter(mqttConfigurationProperties.getSubClientId(), mqttClientFactory);
        // adapter.addTopic("pub/300119110099"); 訂閱主題,也可以放在初始化動態(tài)配置
        adapter.setOutputChannel(mqttInputChannel());
        return adapter;
    }
    /**
     * 配置消息處理器
     *
     * @return
     */
    @Bean
    @ServiceActivator(inputChannel = "mqttInputChannel") // 指定通道
    public MessageHandler messageHandler() {
        return new MqttReceiverMessageHandler();
    }
}

5、消息處理器配置

@Slf4j
@Component
public class MqttReceiverMessageHandler implements MessageHandler {
    @Autowired
    private MqttMessageProcessingService mqttMessageProcessingService;
    @Override
    public void handleMessage(Message<?> message) throws MessagingException {
        MessageHeaders headers = message.getHeaders();
        log.info("線程名稱:{},收到消息,主題:{},消息:{}", Thread.currentThread().getName(), headers.get("mqtt_receivedTopic").toString(), message.getPayload());
        // log.info("收到消息主題:{}", headers.get("mqtt_receivedTopic").toString());
        // log.info("收到消息:{}", message.getPayload());
        // 消息保存到內(nèi)存隊列里面,定時批量入庫,也可以在這里直接入庫
        mqttMessageProcessingService.addMessage(message.getPayload().toString());
    }
}

6、消息主題緩存對象

@Component
public class MqttTopicStore {
    private final ConcurrentHashMap<String, String> topics = new ConcurrentHashMap<>();
    public ConcurrentHashMap<String, String> getTopics() {
        return topics;
    }
}

7、動態(tài)訂閱數(shù)據(jù)庫主題配置

@Slf4j
@Component
public class MqttInit {
    @Autowired
    private MqttPahoMessageDrivenChannelAdapter messageDrivenChannelAdapter;
    @Autowired
    private MqttTopicStore mqttTopicStore;
    @PostConstruct
    public void init() {
        subscribeAllTopics();
    }
    public void subscribeAllTopics() {
        // List<MqttTopicConfig> topics = topicConfigMapper.findAllEnabled();
        // for (MqttTopicConfig topic : topics) {
        //     subscribeTopic(topic);
        // }
        log.info("===================>從數(shù)據(jù)庫里獲取并初始化訂閱所有主題");
        List<String> topics = ListUtil.list(false, "pub/300119110099", "pub1/3010230209810018992", "pub1/30102302098100");
        topics.stream().forEach(t -> {
            messageDrivenChannelAdapter.addTopic(t);
            // 同時往MqttTopicStore.topics中增加一條記錄用于緩存
        });
    }
}

8、消息處理服務

@Service
public class MqttMessageProcessingService {
    @Autowired
    private MqttPahoMessageDrivenChannelAdapter messageDrivenChannelAdapter;
    @Autowired
    private MqttTopicStore mqttTopicStore;
    // 內(nèi)存隊列,用于暫存消息
    private final BlockingQueue<String> messageQueue = new LinkedBlockingQueue<>();
    // 添加消息到隊列
    public void addMessage(String message) {
        messageQueue.add(message);
    }
    /**
     * 可以放到定時任務里面去,注入后取隊列方便維護
     * 定時任務,每5秒執(zhí)行一次 ,建議2分鐘一次 理想的觸發(fā)間隔應略小于數(shù)據(jù)到達間隔,以確保及時處理和插入
     * 如果每 5 分鐘收到一條數(shù)據(jù),可以設置任務執(zhí)行周期為4 分鐘或更短,以便任務有足夠的時間處理數(shù)據(jù),同時減少積壓的可能性。
     */
    @Scheduled(fixedRate = 1 * 60 * 1000)
    public void batchInsertToDatabase() {
        System.out.println("定時任務執(zhí)行中,當前隊列大?。? + messageQueue.size());
        List<String> batch = new ArrayList<>();
        messageQueue.drainTo(batch, 500); // 一次性取最多500條消息
        if (!batch.isEmpty()) {
            // 批量插入數(shù)據(jù)庫
            saveMessagesToDatabase(batch);
        }
    }
    private void saveMessagesToDatabase(List<String> messages) {
        // 假設這是批量插入邏輯
        System.out.println("批量插入數(shù)據(jù)庫,條數(shù):" + messages.size());
        for (String message : messages) {
            System.out.println("插入消息:" + message);
        }
        // 實際數(shù)據(jù)庫操作代碼
    }
    /**
     * 訂閱與取消訂閱定時任務
     */
    public void subscribeAndUnsubscribeTask() {
        // 從數(shù)據(jù)庫獲取所有主題,正常狀態(tài)、刪除狀態(tài)
        // 正常狀態(tài):判斷mqttTopicStore.topics中是否存在,不存在則訂閱,并在mqttTopicStore.topics中增加
        // 刪除狀態(tài): 判斷mqttTopicStore.topics中是否存在,存在則取消訂閱,并在mqttTopicStore.topics中刪除
        // messageDrivenChannelAdapter.addTopic(t);
    }
}

以上是簡單的對接步驟,部分類、方法可以根據(jù)實際情況進行合并處理!?。。?/p>

9、定時任務

@Slf4j
@Configuration
@EnableScheduling
public class MqttJob {
    @Value("${schedule.enable}")
    private boolean enable;
    @Autowired
    private MqttMessageProcessingService mqttMessageProcessingService;
    /**
     * 定時訂閱與取消訂閱主題,從共享主題對象MqttTopicStore里面取出主題列表,然后進行訂閱或取消訂閱
     * 每分鐘一次
     */
    public void subscribeAndUnsubscribe() {
        if (!enable) return;
        mqttMessageProcessingService.subscribeAndUnsubscribeTask();
    }
    /**
     * 定時處理隊列里面的訂閱消息,會有丟失風險,宕機時會丟失隊列里面的消息
     * 每分鐘一次 要考慮一次消息處理的時間;也可先不使用隊列,每次收到消息直接實時入庫,有性能問題時在啟用
     */
    public void batchSaveSubscribeMessage() {
    }
}

到此這篇關于SpringBoot 集成MQTT實現(xiàn)消息訂閱的文章就介紹到這了,更多相關SpringBoot MQTT消息訂閱內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • Springboot 如何實現(xiàn)filter攔截token驗證和跨域

    Springboot 如何實現(xiàn)filter攔截token驗證和跨域

    這篇文章主要介紹了Springboot 如何實現(xiàn)filter攔截token驗證和跨域操作,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-08-08
  • Java加載ICC文件的方法和示例代碼

    Java加載ICC文件的方法和示例代碼

    ICC文件,通常用于顏色管理,定義了如何將一個顏色空間轉(zhuǎn)換為另一個顏色空間,在Java中,我們可能需要加載這些文件來進行顏色轉(zhuǎn)換或管理,本文將為您提供加載ICC文件的方法和示例代碼,需要的朋友參考下吧
    2023-08-08
  • 探究實現(xiàn)Aware接口的原理及使用

    探究實現(xiàn)Aware接口的原理及使用

    這篇文章主要為大家介紹了探究實現(xiàn)Aware接口的原理及使用,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-04-04
  • SpringBoot如何使用自定義注解實現(xiàn)接口限流

    SpringBoot如何使用自定義注解實現(xiàn)接口限流

    這篇文章主要介紹了SpringBoot如何使用自定義注解實現(xiàn)接口限流,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-06-06
  • Freemarker 最簡單的例子程序

    Freemarker 最簡單的例子程序

    Freemarker最簡單的例子程序是通過String來創(chuàng)建模版對象,并執(zhí)行插值處理。
    2016-04-04
  • 通過netty把百度地圖API獲取的地理位置從Android端發(fā)送到Java服務器端的操作方法

    通過netty把百度地圖API獲取的地理位置從Android端發(fā)送到Java服務器端的操作方法

    這篇文章主要介紹了通過netty把百度地圖API獲取的地理位置從Android端發(fā)送到Java服務器端,本文通過示例代碼給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2022-10-10
  • Java 實戰(zhàn)范例之線上新聞平臺系統(tǒng)的實現(xiàn)

    Java 實戰(zhàn)范例之線上新聞平臺系統(tǒng)的實現(xiàn)

    讀萬卷書不如行萬里路,只學書上的理論是遠遠不夠的,只有在實戰(zhàn)中才能獲得能力的提升,本篇文章手把手帶你用java+jsp+jdbc+mysql實現(xiàn)一個線上新聞平臺系統(tǒng),大家可以在過程中查缺補漏,提升水平
    2021-11-11
  • AbstractQueuedSynchronizer內(nèi)部類Node使用講解

    AbstractQueuedSynchronizer內(nèi)部類Node使用講解

    這篇文章主要為大家介紹了AbstractQueuedSynchronizer內(nèi)部類Node使用講解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-07-07
  • QueryWrapper中查詢的坑及解決

    QueryWrapper中查詢的坑及解決

    這篇文章主要介紹了QueryWrapper中查詢的坑及解決方案,具有很好的參考價值,希望對大家有所幫助。
    2022-01-01
  • Mybatis中#{}與${}的區(qū)別詳解

    Mybatis中#{}與${}的區(qū)別詳解

    這篇文章主要介紹了Mybatis中#{}與${}的區(qū)別詳解,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2019-12-12

最新評論

滨海县| 梨树县| 泾阳县| 镇巴县| 甘孜| 怀化市| 垦利县| 泰来县| 都兰县| 濮阳市| 阜新市| 河池市| 温州市| 绥德县| 团风县| 八宿县| 集安市| 海门市| 延吉市| 茶陵县| 镇江市| 齐齐哈尔市| 成都市| 大邑县| 苍梧县| 荣昌县| 辰溪县| 颍上县| 图片| 栾城县| 天峻县| 璧山县| 衡阳市| 顺平县| 泸水县| 理塘县| 桂东县| 磴口县| 菏泽市| 兰溪市| 英吉沙县|