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

RocketMQ中的通信模塊詳解

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

 通信機制

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

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

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

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

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

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

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

Remoting通信類

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

NettyRemotingServer

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

NettyRemotingServer構(gòu)造時主要工作是初始化下列屬性,構(gòu)造時判斷useEpoll來決定EventLoopGroup的實現(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啟動

    @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檢測、連接管理、消息處理
                        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)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

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

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

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

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

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

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

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

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

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

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

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

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

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

    Spring中AOP的切點、通知、切點表達式及知識要點整理

    這篇文章主要介紹了Spring中AOP的切點、通知、切點表達式及知識要點整理,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2023-03-03
  • Java 23種設計模型詳解

    Java 23種設計模型詳解

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

    詳解Spring Boot對 Apache Pulsar的支持

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

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

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

最新評論

永兴县| 清徐县| 景谷| 山丹县| 额济纳旗| 玛沁县| 类乌齐县| 大同县| 疏附县| 清河县| 讷河市| 通榆县| 根河市| 温泉县| 日土县| 河池市| 滁州市| 丰镇市| 仪征市| 屯留县| 宕昌县| 巴楚县| 开封市| 江源县| 凤山市| 常德市| 汾阳市| 雅江县| 怀远县| 普兰店市| 喀喇沁旗| 连平县| 宽城| 扬中市| 浦县| 耿马| 泰来县| 贺州市| 宁夏| 雅江县| 张掖市|