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

RocketMQ中的通信模塊詳解

 更新時(shí)間:2024年01月03日 10:58:31   作者:潛水路人甲  
這篇文章主要介紹了RocketMQ中的通信模塊詳解,RocketMQ消息隊(duì)列集群主要包括NameServer、Broker(Master/Slave)、Producer、Consumer4個(gè)角色,本文我們簡(jiǎn)單來講解一下,需要的朋友可以參考下

 通信機(jī)制

RocketMQ消息隊(duì)列集群主要包括NameServer、Broker(Master/Slave)、Producer、Consumer4個(gè)角色,基本通訊流程如下:

(1) Broker啟動(dòng)后需要完成一次將自己注冊(cè)至NameServer的操作;隨后每隔30s時(shí)間定時(shí)向NameServer上報(bào)Topic路由信息。

(2) 消息生產(chǎn)者Producer作為客戶端發(fā)送消息時(shí)候,需要根據(jù)消息的Topic從本地緩存的TopicPublishInfoTable獲取路由信息。如果沒有則更新路由信息會(huì)從NameServer上重新拉取,同時(shí)Producer會(huì)默認(rèn)每隔30s向NameServer拉取一次路由信息。

(3) 消息生產(chǎn)者Producer根據(jù)2)中獲取的路由信息選擇一個(gè)隊(duì)列(MessageQueue)進(jìn)行消息發(fā)送;Broker作為消息的接收者接收消息并落盤存儲(chǔ)。

(4) 消息消費(fèi)者Consumer根據(jù)2)中獲取的路由信息,并再完成客戶端的負(fù)載均衡后,選擇其中的某一個(gè)或者某幾個(gè)消息隊(duì)列來拉取消息并進(jìn)行消費(fèi)。

從上面1)~3)中可以看出在消息生產(chǎn)者, Broker和NameServer之間都會(huì)發(fā)生通信(這里只說了MQ的部分通信),因此如何設(shè)計(jì)一個(gè)良好的網(wǎng)絡(luò)通信模塊在MQ中至關(guān)重要,它將決定RocketMQ集群整體的消息傳輸能力與最終的性能。

rocketmq-remoting 模塊是 RocketMQ消息隊(duì)列中負(fù)責(zé)網(wǎng)絡(luò)通信的模塊,它幾乎被其他所有需要網(wǎng)絡(luò)通信的模塊(諸如rocketmq-client、rocketmq-broker、rocketmq-namesrv)所依賴和引用。為了實(shí)現(xiàn)客戶端與服務(wù)器之間高效的數(shù)據(jù)請(qǐng)求與接收,RocketMQ消息隊(duì)列自定義了通信協(xié)議并在Netty的基礎(chǔ)之上擴(kuò)展了通信模塊。

Remoting通信類

通信類結(jié)構(gòu):

NettyRemotingServer

NettyRemotingServer為服務(wù)端實(shí)現(xiàn)類,在NamesrvController中被構(gòu)造。

NettyRemotingServer構(gòu)造時(shí)主要工作是初始化下列屬性,構(gòu)造時(shí)判斷useEpoll來決定EventLoopGroup的實(shí)現(xiàn)。

    private final ServerBootstrap serverBootstrap;
    private final EventLoopGroup eventLoopGroupSelector;
    private final EventLoopGroup eventLoopGroupBoss;
    private final NettyServerConfig nettyServerConfig;
 
    private final ExecutorService publicExecutor;
    private final ChannelEventListener channelEventListener;

NettyRemotingServer啟動(dòng)

    @Override
    public void start() {
        this.defaultEventExecutorGroup = new DefaultEventExecutorGroup(
            nettyServerConfig.getServerWorkerThreads(),
            new ThreadFactory() {
 
                private AtomicInteger threadIndex = new AtomicInteger(0);
 
                @Override
                public Thread newThread(Runnable r) {
                    return new Thread(r, "NettyServerCodecThread_" + this.threadIndex.incrementAndGet());
                }
            });
 
        prepareSharableHandlers();
 
        ServerBootstrap childHandler =
            this.serverBootstrap.group(this.eventLoopGroupBoss, this.eventLoopGroupSelector)
                .channel(useEpoll() ? EpollServerSocketChannel.class : NioServerSocketChannel.class)
                .option(ChannelOption.SO_BACKLOG, 1024)
                .option(ChannelOption.SO_REUSEADDR, true)
                .option(ChannelOption.SO_KEEPALIVE, false)
                .childOption(ChannelOption.TCP_NODELAY, true)
                .childOption(ChannelOption.SO_SNDBUF, nettyServerConfig.getServerSocketSndBufSize())
                .childOption(ChannelOption.SO_RCVBUF, nettyServerConfig.getServerSocketRcvBufSize())
                .localAddress(new InetSocketAddress(this.nettyServerConfig.getListenPort()))
                .childHandler(new ChannelInitializer<SocketChannel>() {
                    @Override
                    public void initChannel(SocketChannel ch) throws Exception {
					//添加handler,握手、編解碼、idle檢測(cè)、連接管理、消息處理
                        ch.pipeline()
                            .addLast(defaultEventExecutorGroup, HANDSHAKE_HANDLER_NAME, handshakeHandler)
                            .addLast(defaultEventExecutorGroup,
                                encoder,
                                new NettyDecoder(),
                                new IdleStateHandler(0, 0, nettyServerConfig.getServerChannelMaxIdleTimeSeconds()),
                                connectionManageHandler,
                                serverHandler
                            );
                    }
                });
 
        if (nettyServerConfig.isServerPooledByteBufAllocatorEnable()) {
            childHandler.childOption(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT);
        }
 
        try {
            ChannelFuture sync = this.serverBootstrap.bind().sync();
            InetSocketAddress addr = (InetSocketAddress) sync.channel().localAddress();
            this.port = addr.getPort();
        } catch (InterruptedException e1) {
            throw new RuntimeException("this.serverBootstrap.bind().sync() InterruptedException", e1);
        }
 
        if (this.channelEventListener != null) {
            this.nettyEventExecutor.start();
        }
 
        this.timer.scheduleAtFixedRate(new TimerTask() {
 
            @Override
            public void run() {
                try {
                    NettyRemotingServer.this.scanResponseTable();
                } catch (Throwable e) {
                    log.error("scanResponseTable exception", e);
                }
            }
        }, 1000 * 3, 1000);
    }

主要關(guān)注initChannel方法中添加的handler,NettyConnectManageHandler  

    //消息處理核心類   
    @ChannelHandler.Sharable
    class NettyServerHandler extends SimpleChannelInboundHandler<RemotingCommand> {
        @Override
        protected void channelRead0(ChannelHandlerContext ctx, RemotingCommand msg) throws Exception {
            processMessageReceived(ctx, msg);
        }
    }
    //連接管理處理類
    @ChannelHandler.Sharable
    class NettyConnectManageHandler extends ChannelDuplexHandler {
        @Override
        public void channelActive(ChannelHandlerContext ctx) throws Exception {
            final String remoteAddress = RemotingHelper.parseChannelRemoteAddr(ctx.channel());
            log.info("NETTY SERVER PIPELINE: channelActive, the channel[{}]", remoteAddress);
            super.channelActive(ctx);
            if (NettyRemotingServer.this.channelEventListener != null) {
                NettyRemotingServer.this.putNettyEvent(new NettyEvent(NettyEventType.CONNECT, remoteAddress, ctx.channel()));
            }
        }
        @Override
        public void channelInactive(ChannelHandlerContext ctx) throws Exception {
            final String remoteAddress = RemotingHelper.parseChannelRemoteAddr(ctx.channel());
            log.info("NETTY SERVER PIPELINE: channelInactive, the channel[{}]", remoteAddress);
            super.channelInactive(ctx);
            if (NettyRemotingServer.this.channelEventListener != null) {
                NettyRemotingServer.this.putNettyEvent(new NettyEvent(NettyEventType.CLOSE, remoteAddress, ctx.channel()));
            }
        }
        @Override
        public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception {
            if (evt instanceof IdleStateEvent) {
                IdleStateEvent event = (IdleStateEvent) evt;
                if (event.state().equals(IdleState.ALL_IDLE)) {
                    final String remoteAddress = RemotingHelper.parseChannelRemoteAddr(ctx.channel());
                    log.warn("NETTY SERVER PIPELINE: IDLE exception [{}]", remoteAddress);
                    RemotingUtil.closeChannel(ctx.channel());
                    if (NettyRemotingServer.this.channelEventListener != null) {
                        NettyRemotingServer.this
                            .putNettyEvent(new NettyEvent(NettyEventType.IDLE, remoteAddress, ctx.channel()));
                    }
                }
            }
            ctx.fireUserEventTriggered(evt);
        }
        @Override
        public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
            final String remoteAddress = RemotingHelper.parseChannelRemoteAddr(ctx.channel());
            log.warn("NETTY SERVER PIPELINE: exceptionCaught {}", remoteAddress);
            log.warn("NETTY SERVER PIPELINE: exceptionCaught exception.", cause);
            if (NettyRemotingServer.this.channelEventListener != null) {
                NettyRemotingServer.this.putNettyEvent(new NettyEvent(NettyEventType.EXCEPTION, remoteAddress, ctx.channel()));
            }
            RemotingUtil.closeChannel(ctx.channel());
        }
    }

到此這篇關(guān)于RocketMQ中的通信模塊詳解的文章就介紹到這了,更多相關(guān)RocketMQ通信模塊內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java基于UDP協(xié)議的聊天室功能

    Java基于UDP協(xié)議的聊天室功能

    這篇文章主要為大家詳細(xì)介紹了Java基于UDP協(xié)議的聊天室功能,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2022-09-09
  • idea日志亂碼和tomcat日志亂碼問題的解決方法

    idea日志亂碼和tomcat日志亂碼問題的解決方法

    這篇文章主要介紹了idea日志亂碼和tomcat日志亂碼問題的解決方法,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2020-08-08
  • ShardingProxy讀寫分離之原理、配置與實(shí)踐過程

    ShardingProxy讀寫分離之原理、配置與實(shí)踐過程

    ShardingProxy是Apache?ShardingSphere的數(shù)據(jù)庫中間件,通過三層架構(gòu)實(shí)現(xiàn)讀寫分離,解決高并發(fā)場(chǎng)景下數(shù)據(jù)庫性能瓶頸,其核心功能包括SQL路由、負(fù)載均衡、數(shù)據(jù)一致性保障和故障轉(zhuǎn)移,支持主從架構(gòu)下的透明分庫分表及讀寫分流,廣泛應(yīng)用于微服務(wù)和高流量業(yè)務(wù)系統(tǒng)
    2025-08-08
  • SpringBoot的ConfigurationProperties或Value注解無效問題及解決

    SpringBoot的ConfigurationProperties或Value注解無效問題及解決

    在SpringBoot項(xiàng)目開發(fā)中,全局靜態(tài)配置類讀取application.yml或application.properties文件時(shí),可能會(huì)遇到配置值始終為null的問題,這通常是因?yàn)樵趧?chuàng)建靜態(tài)屬性后,IDE自動(dòng)生成的Get/Set方法包含了static關(guān)鍵字
    2024-11-11
  • 配置Spring4.0注解Cache+Redis緩存的用法

    配置Spring4.0注解Cache+Redis緩存的用法

    本篇文章主要介紹了詳解配置Spring4.0注解Cache+Redis緩存的用法,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2017-04-04
  • SpringBoot實(shí)現(xiàn)文件上傳下載的7種方法

    SpringBoot實(shí)現(xiàn)文件上傳下載的7種方法

    文件上傳下載功能是Web應(yīng)用中的常見需求,從簡(jiǎn)單的用戶頭像上傳到大型文件的傳輸與共享,都需要可靠的文件處理機(jī)制,下面我們來看看SpringBoot中處理文件上傳下載的7種方法吧
    2025-05-05
  • Spring中AOP的切點(diǎn)、通知、切點(diǎn)表達(dá)式及知識(shí)要點(diǎn)整理

    Spring中AOP的切點(diǎn)、通知、切點(diǎn)表達(dá)式及知識(shí)要點(diǎn)整理

    這篇文章主要介紹了Spring中AOP的切點(diǎn)、通知、切點(diǎn)表達(dá)式及知識(shí)要點(diǎn)整理,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2023-03-03
  • Java 23種設(shè)計(jì)模型詳解

    Java 23種設(shè)計(jì)模型詳解

    本文主要介紹Java 23種設(shè)計(jì)模型,這里整理了詳細(xì)的資料,及實(shí)現(xiàn)各種設(shè)計(jì)模型的示例代碼,有需要的小伙伴可以參考下
    2016-09-09
  • 詳解Spring Boot對(duì) Apache Pulsar的支持

    詳解Spring Boot對(duì) Apache Pulsar的支持

    Spring Boot通過提供spring-pulsar和spring-pulsar-reactive自動(dòng)配置支持Apache Pulsar,類路徑中這些依賴存在時(shí),Spring Boot自動(dòng)配置命令式和反應(yīng)式Pulsar組件,PulsarClient自動(dòng)注冊(cè),默認(rèn)連接本地Pulsar實(shí)例,感興趣的朋友一起看看吧
    2024-11-11
  • 詳解Java的四種引用方式及其區(qū)別

    詳解Java的四種引用方式及其區(qū)別

    這篇文章主要介紹了Java的四種引用方式 ,主要主要包括強(qiáng)引用,軟引用,弱引用,虛引用,稍微整理精簡(jiǎn)一下做下分享,具有一定的參考價(jià)值,需要的朋友可以參考下
    2018-12-12

最新評(píng)論

安远县| 分宜县| 满洲里市| 大新县| 定边县| 朝阳县| 鄂州市| 加查县| 尚志市| 尚义县| 龙山县| 新余市| 岳阳市| 高淳县| 锦州市| 宜兰市| 临澧县| 双牌县| 瑞金市| 丰顺县| 通渭县| 武汉市| 镇康县| 正蓝旗| 安乡县| 蕲春县| 松桃| 江山市| 渝北区| 西丰县| 宣武区| 兴山县| 镶黄旗| 天气| 兰西县| 勃利县| 剑河县| 日土县| 弥渡县| 永寿县| 海南省|