zoukankan      html  css  js  c++  java
  • Rocketmq消息持久化

    本文编写,参考:https://my.oschina.net/bieber/blog/725646

    producer Send()的Message最终将由broker处理,处理类为:SendMessageProcessor ,处理方法:processRequet.

    public class SendMessageProcessor extends AbstractSendMessageProcessor implements NettyRequestProcessor {

    private List<ConsumeMessageHook> consumeMessageHookList;

    public SendMessageProcessor(final BrokerController brokerController) {
    super(brokerController);
    }

    @Override
    public RemotingCommand processRequest(ChannelHandlerContext ctx, RemotingCommand request) throws RemotingCommandException {}
    上述方法,并不是直接处理消息,而是交由MessageStore处理,相关代码如下:
    MessageExtBrokerInner msgInner = new MessageExtBrokerInner();
    msgInner.setTopic(requestHeader.getTopic());
    msgInner.setQueueId(queueIdInt);
    //......
    PutMessageResult putMessageResult = this.brokerController.getMessageStore().putMessage(msgInner);

    然而MessageStore也不直接持久化消息,转交给 CommitLog
    long beginTime = this.getSystemClock().now();
    PutMessageResult result = this.commitLog.putMessages(messageExtBatch);

    从MappedFileQueue中取出最新的一条:
    MappedFile mappedFile = this.mappedFileQueue.getLastMappedFile();
    //写消息
    result = mappedFile.appendMessages(messageExtBatch, this.appendMessageCallback);
    //持久化到磁盘,最终通过FileChannel持久化到文件
    handleDiskFlush(result, putMessageResult, messageExtBatch);

    handleHA(result, putMessageResult, messageExtBatch);


    2.cousumer 从broker读消息。
    消费者从broker读取消息经由PullMessageProcessor类处理的,processRequest()方法处理请求:
    RemotingCommand processRequest(final Channel channel, RemotingCommand request, boolean brokerAllowSuspend)

    经过一系列的判断处理,之后交由 MessageStore:
    final GetMessageResult getMessageResult =
    this.brokerController.getMessageStore().getMessage(requestHeader.getConsumerGroup(), requestHeader.getTopic(),
    requestHeader.getQueueId(), requestHeader.getQueueOffset(), requestHeader.getMaxMsgNums(), messageFilter);
    读取消息。
    之后交由commitLog,读出消息,
    SelectMappedBufferResult selectResult = this.commitLog.getMessage(offsetPy, sizePy);
    可以看到是先从ConsumerQueue中获取消息索引,然后再从commitlog中读取消息内容。这些内容也是在存储消息的时候写入的。
    相关也可参考:http://jm.taobao.org/2017/01/12/rocketmq-quick-start-in-10-minutes/



  • 相关阅读:
    内存映射文件原理探索(转)
    inux内存映射和共享内存理解和区别
    MySQL中的sleep函数介绍
    flask源码解析之上下文为什么用栈
    linux system函数详解
    Python中的可迭代对象、迭代器和生成器,协程的异同点
    GB2312汉字区位码、交换码和机内码转换方法 (ZT)
    pthread_cond_signal该放在什么地方?
    IPC介绍——10个ipcs例子
    lsof
  • 原文地址:https://www.cnblogs.com/itdev/p/7086322.html
Copyright © 2011-2022 走看看