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

SpringBoot整合Netty開發(fā)MQTT服務端

 更新時間:2025年07月22日 09:45:23   作者:大魚>  
Netty是一款基于NIO(Nonblocking I/O,非阻塞IO)開發(fā)的網絡通信框架,本文主要介紹了SpringBoot如何整合Netty開發(fā)MQTT服務端,感興趣的小伙伴可以了解下

Netty認知

Netty是一款基于NIO(Nonblocking I/O,非阻塞IO)開發(fā)的網絡通信框架,相比傳統(tǒng)Socket,在并發(fā)性方面有著很大的提升。關于NIO,BIO,AIO之間的區(qū)別,可以參考這篇博客:Java中AIO、BIO、NIO應用場景及區(qū)別 

MQTT服務端實現(xiàn)

首先我們啟動一個tcp服務,這里我用到了Redis與RabbitMQ,主要是與分布式WEB平臺之間好對接

@Component
public class ApplicationEventListener implements CommandLineRunner {
    @Value("${spring.application.name}")
    private String nodeName;

    @Value("${gnss.mqttserver.tcpPort}")
    private int tcpPort;

    @Override
    public void run(String... args) throws Exception {
        //啟動TCP服務
        startTcpServer();

        //清除Redis所有此節(jié)點的在線終端
        RedisService redisService = SpringBeanService.getBean(RedisService.class);
        redisService.deleteAllOnlineTerminals(nodeName);

        //將所有此節(jié)點的終端設置為離線
        RabbitMessageSender messageSender = SpringBeanService.getBean(RabbitMessageSender.class);
        messageSender.noticeAllOffline(nodeName);
    }

    /**
     * 啟動TCP服務
     *
     * @throws Exception
     */
    private void startTcpServer() throws Exception {
        //計數器,必須等到所有服務啟動成功才能進行后續(xù)的操作
        final CountDownLatch countDownLatch = new CountDownLatch(1);
        //啟動TCP服務
        TcpServer tcpServer = new TcpServer(tcpPort, ProtocolEnum.MqttCommon, countDownLatch);
        tcpServer.start();
        //等待啟動完成
        countDownLatch.await();
    }
}

接下來我們編寫一個TcpServer類實現(xiàn)TCP服務

@Slf4j
public class TcpServer extends Thread{
    private int port;

    private ProtocolEnum protocolType;

    private EventLoopGroup bossGroup;

    private EventLoopGroup workerGroup;

    private ServerBootstrap serverBootstrap = new ServerBootstrap();

    private CountDownLatch countDownLatch;

    public TcpServer(int port, ProtocolEnum protocolType, CountDownLatch countDownLatch) {
        this.port = port;
        this.protocolType = protocolType;
        this.countDownLatch = countDownLatch;

        bossGroup = new NioEventLoopGroup(1);
        workerGroup = SpringBeanService.getBean("workerGroup", EventLoopGroup.class);
        final EventExecutorGroup executorGroup = SpringBeanService.getBean("executorGroup", EventExecutorGroup.class);
        serverBootstrap.group(bossGroup, workerGroup)
                .channel(NioServerSocketChannel.class)
                .option(ChannelOption.SO_BACKLOG, 1024)
                .childOption(ChannelOption.SO_KEEPALIVE, true)
                .childOption(ChannelOption.TCP_NODELAY, true)
                .childHandler(new ChannelInitializer<SocketChannel>() {

                    @Override
                    protected void initChannel(SocketChannel ch) throws Exception {
                        ch.pipeline().addLast(new IdleStateHandler(MqttConstant.READER_IDLE_TIME, 0, 0, TimeUnit.SECONDS));
                        ch.pipeline().addLast("encoder", MqttEncoder.INSTANCE);
                        ch.pipeline().addLast("decoder", new MqttDecoder());
                        ch.pipeline().addLast(executorGroup, MqttBusinessHandler.INSTANCE);
                    }
                });
    }

    @Override
    public void run() {
        bind();
    }

    /**
     * 綁定端口啟動服務
     */
    private void bind() {
        serverBootstrap.bind(port).addListener(future -> {
            if (future.isSuccess()) {
                log.info("{} MQTT服務器啟動,端口:{}", protocolType, port);
                countDownLatch.countDown();
            } else {
                log.error("{} MQTT服務器啟動失敗,端口:{}", protocolType, port, future.cause());
                System.exit(-1);
            }
        });
    }

    /**
     * 關閉服務端
     */
    public void shutdown() {
        workerGroup.shutdownGracefully();
        bossGroup.shutdownGracefully();
        log.info("{} TCP服務器關閉,端口:{}", protocolType, port);
    }
}

編寫一個解碼器MqttBusinessHandler,實現(xiàn)對MQTT消息接收與處理

@Slf4j
@ChannelHandler.Sharable
public class MqttBusinessHandler extends SimpleChannelInboundHandler<Object> {
    public static final MqttBusinessHandler INSTANCE = new MqttBusinessHandler();
    private MqttMsgBack mqttMsgBack;
    private MqttBusinessHandler() {
        mqttMsgBack= MqttMsgBack.INSTANCE;
    }

    /**
     * 接收到消息后處理
     * @param ctx
     * @param msg
     * @throws Exception
     */
    @Override
    protected void channelRead0(ChannelHandlerContext ctx, Object msg) throws Exception {
        if (null != msg) {
            MqttMessage mqttMessage = (MqttMessage) msg;
            MqttFixedHeader mqttFixedHeader = mqttMessage.fixedHeader();
            Channel channel = ctx.channel();
            if(mqttFixedHeader.messageType().equals(MqttMessageType.CONNECT)){
                //在一個網絡連接上,客戶端只能發(fā)送一次CONNECT報文。服務端必須將客戶端發(fā)送的第二個CONNECT報文當作協(xié)議違規(guī)處理并斷開客戶端的連接
                //建議connect消息單獨處理,用來對客戶端進行認證管理等 這里直接返回一個CONNACK消息
                mqttMsgBack.connectionAck(ctx, mqttMessage);
            }

            switch (mqttFixedHeader.messageType()){
                //客戶端發(fā)布消息
                case PUBLISH:
                    mqttMsgBack.publishAck(ctx, mqttMessage);
                    break;
                //發(fā)布釋放
                case PUBREL:
                    mqttMsgBack.publishComp(ctx, mqttMessage);
                    break;
                //訂閱主題
                case SUBSCRIBE:
                    mqttMsgBack.subscribeAck(ctx, mqttMessage);
                    break;
                //取消訂閱主題
                case UNSUBSCRIBE:
                    mqttMsgBack.unsubscribeAck(ctx, mqttMessage);
                    break;
                //客戶端發(fā)送心跳報文
                case PINGREQ:
                    mqttMsgBack.pingResp(ctx, mqttMessage);
                    break;
                //客戶端主動斷開連接
                case DISCONNECT:
                    break;
                default:
                    break;
            }
        }
    }

    @Override
    public void channelInactive(ChannelHandlerContext ctx) throws Exception {
        log.info("終端關閉連接,IP信息:{}", CommonUtil.getClientAddress(ctx));
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
        ctx.close();
        log.error("終端連接異常,IP信息:{}", CommonUtil.getClientAddress(ctx), cause);
    }

    /**
     * 	服務端當讀超時時會調用這個方法
     */
    @Override
    public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception, IOException {
        ctx.close();
        log.error("讀超時,IP信息:{}", CommonUtil.getClientAddress(ctx), evt);
    }

我們對接收到的消息進行業(yè)務處理

@Slf4j
public class MqttMsgBack {
    public static final MqttMsgBack INSTANCE = new MqttMsgBack();
    private RedisService redisService;
    private RabbitMessageSender messageSender;
    private Environment environment;
    private MessageServiceProvider messageServiceProvider;

    private MqttMsgBack() {
        redisService = SpringBeanService.getBean(RedisService.class);
        messageSender = SpringBeanService.getBean(RabbitMessageSender.class);
        environment = SpringBeanService.getBean(Environment.class);
        messageServiceProvider = SpringBeanService.getBean(MessageServiceProvider.class);
    }

    /**
     * 	確認連接請求
     * @param ctx
     * @param mqttMessage
     */
    public void connectionAck (ChannelHandlerContext ctx, MqttMessage mqttMessage) {
        MqttConnectMessage mqttConnectMessage = (MqttConnectMessage) mqttMessage;
        MqttFixedHeader mqttFixedHeaderInfo = mqttConnectMessage.fixedHeader();
        MqttConnectVariableHeader mqttConnectVariableHeaderInfo = mqttConnectMessage.variableHeader();
        //構建返回報文, 可變報頭
        MqttConnAckVariableHeader mqttConnAckVariableHeaderBack = new MqttConnAckVariableHeader(MqttConnectReturnCode.CONNECTION_ACCEPTED, mqttConnectVariableHeaderInfo.isCleanSession());
        //構建返回報文, 固定報頭
        MqttFixedHeader mqttFixedHeaderBack = new MqttFixedHeader(MqttMessageType.CONNACK,mqttFixedHeaderInfo.isDup(), MqttQoS.AT_MOST_ONCE, mqttFixedHeaderInfo.isRetain(), 0x02);
        //構建連接回復消息體
        MqttConnAckMessage connAck = new MqttConnAckMessage(mqttFixedHeaderBack, mqttConnAckVariableHeaderBack);
        ctx.writeAndFlush(connAck);
        //獲取連接者的ClientId
        String clientIdentifier = mqttConnectMessage.payload().clientIdentifier();
        //查詢終端號碼有無在平臺注冊
        TerminalProto terminalInfo = redisService.getTerminalInfoByTerminalNum(clientIdentifier);
        if (terminalInfo == null) {
            log.error("終端登錄失敗,未找到終端信息,終端號:{},IP信息:{}", clientIdentifier, CommonUtil.getClientAddress(ctx));
            ctx.close();
            return;
        }
        //設置節(jié)點名
        terminalInfo.setNodeName(environment.getProperty("spring.application.name"));
        //保存終端信息和消息流水號到上下文屬性中
        Session session = new Session(terminalInfo);
        ChannelHandlerContext oldCtx = SessionUtil.bindSession(session, ctx);
        if (oldCtx == null) {
            log.info("終端登錄成功,終端ID:{},終端號:{},IP信息:{}", terminalInfo.getTerminalStrId(), clientIdentifier, CommonUtil.getClientAddress(ctx));
        } else {
            log.info("終端重復登錄關閉上一個連接,終端ID:{},終端號:{},IP信息:{}", terminalInfo.getTerminalStrId(), clientIdentifier, CommonUtil.getClientAddress(ctx));
            oldCtx.close();
        }
        //通知上線
        messageSender.noticeOnline(terminalInfo);
        log.info("終端登錄成功,終端號:{},IP信息:{}", clientIdentifier, CommonUtil.getClientAddress(ctx));
    }

    /**
     * 	根據qos發(fā)布確認
     * @param ctx
     * @param mqttMessage
     */
    public void publishAck (ChannelHandlerContext ctx, MqttMessage mqttMessage) {
        MqttPublishMessage mqttPublishMessage = (MqttPublishMessage) mqttMessage;
        MqttFixedHeader mqttFixedHeaderInfo = mqttPublishMessage.fixedHeader();
        MqttQoS qos = (MqttQoS) mqttFixedHeaderInfo.qosLevel();
        //得到主題
        String topicName = mqttPublishMessage.variableHeader().topicName();
        //獲取消息體
        ByteBuf msgBodyBuf = mqttPublishMessage.payload();
        log.info("收到:{}", ByteBufUtil.hexDump(msgBodyBuf));
        MqttCommonMessage msg=new MqttCommonMessage();
        msg.setTerminalNum(SessionUtil.getTerminalInfo(ctx).getTerminalNum());
        msg.setStrMsgId(topicName);
        //根據主題獲取對應的主題消息處理器
        BaseMessageService messageService = messageServiceProvider.getMessageService(topicName);
        try {
            Object result = messageService.process(ctx, msg, msgBodyBuf);
            log.info("收到{}({}),終端ID:{},內容:{}", messageService.getDesc(), topicName,msg.getTerminalNum(), msg.getMsgBodyItems());
        } catch (Exception e) {
            log.error("收到{}({}),消息異常,終端ID:{},消息體:{}", messageService.getDesc(), topicName,msg.getTerminalNum(),ByteBufUtil.hexDump(msgBodyBuf), e);
        }
        switch (qos) {
            //至多一次
            case AT_MOST_ONCE:
                break;
            //至少一次
            case AT_LEAST_ONCE:
                //構建返回報文, 可變報頭
                MqttMessageIdVariableHeader mqttMessageIdVariableHeaderBack = MqttMessageIdVariableHeader.from(mqttPublishMessage.variableHeader().packetId());
                //構建返回報文, 固定報頭
                MqttFixedHeader mqttFixedHeaderBack = new MqttFixedHeader(MqttMessageType.PUBACK,mqttFixedHeaderInfo.isDup(), MqttQoS.AT_MOST_ONCE, mqttFixedHeaderInfo.isRetain(), 0x02);
                //構建PUBACK消息體
                MqttPubAckMessage pubAck = new MqttPubAckMessage(mqttFixedHeaderBack, mqttMessageIdVariableHeaderBack);
                log.info("Qos:AT_LEAST_ONCE:{}",pubAck.toString());
                ctx.writeAndFlush(pubAck);
                break;
            //剛好一次
            case EXACTLY_ONCE:
                //構建返回報文,固定報頭
                MqttFixedHeader mqttFixedHeaderBack2 = new MqttFixedHeader(MqttMessageType.PUBREC,false, MqttQoS.AT_LEAST_ONCE,false,0x02);
                //構建返回報文,可變報頭
                MqttMessageIdVariableHeader mqttMessageIdVariableHeaderBack2 = MqttMessageIdVariableHeader.from(mqttPublishMessage.variableHeader().packetId());
                MqttMessage mqttMessageBack = new MqttMessage(mqttFixedHeaderBack2,mqttMessageIdVariableHeaderBack2);
                log.info("Qos:EXACTLY_ONCE回復:{}"+mqttMessageBack.toString());
                ctx.writeAndFlush(mqttMessageBack);
                break;
            default:
                break;
        }
    }

    /**
     * 發(fā)布完成 qos2
     * @param ctx
     * @param mqttMessage
     */
    public void publishComp (ChannelHandlerContext ctx, MqttMessage mqttMessage) {

        MqttMessageIdVariableHeader messageIdVariableHeader = (MqttMessageIdVariableHeader) mqttMessage.variableHeader();
        //構建返回報文, 固定報頭
        MqttFixedHeader mqttFixedHeaderBack = new MqttFixedHeader(MqttMessageType.PUBCOMP,false, MqttQoS.AT_MOST_ONCE,false,0x02);
        //構建返回報文, 可變報頭
        MqttMessageIdVariableHeader mqttMessageIdVariableHeaderBack = MqttMessageIdVariableHeader.from(messageIdVariableHeader.messageId());
        MqttMessage mqttMessageBack = new MqttMessage(mqttFixedHeaderBack,mqttMessageIdVariableHeaderBack);
        log.info("發(fā)布完成回復:{}"+mqttMessageBack.toString());
        ctx.writeAndFlush(mqttMessageBack);
    }

    /**
     * 	訂閱確認
     * @param ctx
     * @param mqttMessage
     */
    public void subscribeAck(ChannelHandlerContext ctx, MqttMessage mqttMessage) {
        MqttSubscribeMessage mqttSubscribeMessage = (MqttSubscribeMessage) mqttMessage;
        MqttMessageIdVariableHeader messageIdVariableHeader = mqttSubscribeMessage.variableHeader();
        //構建返回報文, 可變報頭
        MqttMessageIdVariableHeader variableHeaderBack = MqttMessageIdVariableHeader.from(messageIdVariableHeader.messageId());
        Set<String> topics = mqttSubscribeMessage.payload().topicSubscriptions().stream().map(mqttTopicSubscription -> mqttTopicSubscription.topicName()).collect(Collectors.toSet());
        List<Integer> grantedQoSLevels = new ArrayList<>(topics.size());
        for (int i = 0; i < topics.size(); i++) {
            grantedQoSLevels.add(mqttSubscribeMessage.payload().topicSubscriptions().get(i).qualityOfService().value());
        }
        //	構建返回報文	有效負載
        MqttSubAckPayload payloadBack = new MqttSubAckPayload(grantedQoSLevels);
        //	構建返回報文	固定報頭
        MqttFixedHeader mqttFixedHeaderBack = new MqttFixedHeader(MqttMessageType.SUBACK, false, MqttQoS.AT_MOST_ONCE, false, 2+topics.size());
        //	構建返回報文	訂閱確認
        MqttSubAckMessage subAck = new MqttSubAckMessage(mqttFixedHeaderBack,variableHeaderBack, payloadBack);
        log.info("訂閱回復:{}", subAck.toString());
        ctx.writeAndFlush(subAck);
    }

    /**
     * 取消訂閱確認
     * @param ctx
     * @param mqttMessage
     */
    public void unsubscribeAck(ChannelHandlerContext ctx, MqttMessage mqttMessage) {
        MqttMessageIdVariableHeader messageIdVariableHeader = (MqttMessageIdVariableHeader) mqttMessage.variableHeader();
        //	構建返回報文	可變報頭
        MqttMessageIdVariableHeader variableHeaderBack = MqttMessageIdVariableHeader.from(messageIdVariableHeader.messageId());
        //	構建返回報文	固定報頭
        MqttFixedHeader mqttFixedHeaderBack = new MqttFixedHeader(MqttMessageType.UNSUBACK, false, MqttQoS.AT_MOST_ONCE, false, 2);
        //	構建返回報文	取消訂閱確認
        MqttUnsubAckMessage unSubAck = new MqttUnsubAckMessage(mqttFixedHeaderBack,variableHeaderBack);
        log.info("取消訂閱回復:{}",unSubAck.toString());
        ctx.writeAndFlush(unSubAck);
    }

    /**
     * 心跳響應
     * @param ctx
     * @param mqttMessage
     */
    public void pingResp (ChannelHandlerContext ctx, MqttMessage mqttMessage) {
        MqttFixedHeader fixedHeader = new MqttFixedHeader(MqttMessageType.PINGRESP, false, MqttQoS.AT_MOST_ONCE, false, 0);
        MqttMessage mqttMessageBack = new MqttMessage(fixedHeader);
        log.info("心跳回復:{}", mqttMessageBack.toString());
        ctx.writeAndFlush(mqttMessageBack);
    }
}

我們可以根據客戶端發(fā)布消息的主題匹配不同的處理器

 最后,我們在對應的處理器里面實現(xiàn)對主題消息的處理邏輯,比如:定位消息,指令消息等等,比如簡單實現(xiàn)對定位數據Location主題的消息處理

@Slf4j
@MessageService(strMessageId = "Location", desc = "定位")
public class LocationMessageService extends BaseMessageService<MqttCommonMessage> {
    @Autowired
    private RabbitMessageSender messageSender;

    @Override
    public Object process(ChannelHandlerContext ctx, MqttCommonMessage msg, ByteBuf msgBodyBuf) throws Exception {
        byte[] msgByteArr = new byte[msgBodyBuf.readableBytes()];
        msgBodyBuf.readBytes(msgByteArr);
        String data = new String(msgByteArr);
        msg.putMessageBodyItem("位置", data);
        return null;
    }
}

后續(xù)

目前僅僅是實現(xiàn)MQTT服務端消息接收與消息回復,后續(xù)可以根據接入的物聯(lián)網設備進行對應主題消息的業(yè)務處理

到此這篇關于SpringBoot整合Netty開發(fā)MQTT服務端的文章就介紹到這了,更多相關SpringBoot MQTT服務端內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • JDK1.8使用的垃圾回收器和執(zhí)行GC的時長以及GC的頻率方式

    JDK1.8使用的垃圾回收器和執(zhí)行GC的時長以及GC的頻率方式

    這篇文章主要介紹了JDK1.8使用的垃圾回收器和執(zhí)行GC的時長以及GC的頻率方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-05-05
  • Java編程一道多線程問題實例代碼

    Java編程一道多線程問題實例代碼

    這篇文章主要介紹了Java編程一道多線程問題實例代碼,分享了相關代碼示例,小編覺得還是挺不錯的,具有一定借鑒價值,需要的朋友可以參考下
    2018-02-02
  • java實現(xiàn)的xml格式化實現(xiàn)代碼

    java實現(xiàn)的xml格式化實現(xiàn)代碼

    這篇文章主要介紹了java實現(xiàn)的xml格式化實現(xiàn)代碼,需要的朋友可以參考下
    2016-11-11
  • IDEA全局查找關鍵字的用法解讀

    IDEA全局查找關鍵字的用法解讀

    這篇文章主要介紹了IDEA全局查找關鍵字的用法解讀,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-02-02
  • 詳解Java類動態(tài)加載和熱替換

    詳解Java類動態(tài)加載和熱替換

    本文主要介紹類加載器、自定義類加載器及類的加載和卸載等內容,并舉例介紹了Java類的熱替換。
    2021-05-05
  • 如何設置springboot啟動端口

    如何設置springboot啟動端口

    spring boot是個好東西,可以不用容器直接在main方法中啟動,而且無需配置文件,方便快速搭建環(huán)境。下面給大家介紹springboot啟動端口的設置方法和spring boot創(chuàng)建應用端口沖突8080 問題,感興趣的朋友一起看看吧
    2017-08-08
  • MyBatis實現(xiàn)三級樹查詢的示例代碼

    MyBatis實現(xiàn)三級樹查詢的示例代碼

    在實際項目開發(fā)中,樹形結構的數據查詢是一個非常常見的需求,比如組織架構、菜單管理、地區(qū)選擇等場景都需要處理樹形數據,本文將詳細講解如何使用MyBatis實現(xiàn)三級樹形數據的查詢,需要的朋友可以參考下
    2024-12-12
  • Java虛擬線程(VirtualThread)使用詳解

    Java虛擬線程(VirtualThread)使用詳解

    這篇文章主要介紹了Java虛擬線程(VirtualThread)使用,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2025-05-05
  • Spring 切面執(zhí)行鏈的實現(xiàn)示例

    Spring 切面執(zhí)行鏈的實現(xiàn)示例

    在實際應用中,一個方法通常會被多個切面攔截,本文主要介紹了Spring 切面執(zhí)行鏈的實現(xiàn)示例,具有一定的參考價值,感興趣的可以了解一下
    2025-10-10
  • SpringBoot使用Jackson詳解

    SpringBoot使用Jackson詳解

    Spring?Boot中使用Jackson處理JavaBean序列化為JSON格式,常用框架包括Jackson、Fastjson和Gson,Jackson是Spring?Boot默認的JSON處理庫,常用注解如@JsonProperty、@JsonIgnore、@JsonFormat等,用于自定義序列化和反序列化行為
    2025-02-02

最新評論

沂水县| 广东省| 甘谷县| 扶风县| 宁阳县| 杨浦区| 湘潭市| 廊坊市| 玉树县| 东兰县| 宜黄县| 河池市| 广元市| 白山市| 凌源市| 洛宁县| 温州市| 彭阳县| 吉林省| 南丹县| 杭锦旗| 云梦县| 乳山市| 辰溪县| 林口县| 黄冈市| 绿春县| 九台市| 安化县| 赫章县| 应用必备| 崇左市| 安陆市| 阳东县| 庄浪县| 方正县| 凉城县| 九台市| 凤城市| 玉溪市| 华亭县|