本文分析broker响应producer发送非延迟消息的整个过程。
分析的qmq版本是1.1.3.5。
重点结合源码分析如下几个问题:
broker接收消息如何做到全异步处理?
broker是批量接收消息的,那如何做到每个msg独立处理的同时又合并返回结果?
回写此次请求的时间点。
那如何做到的高可用?
首先看下broker接收消息到写入文件的整个过程,如下:
结合上图,从两方面切入分析:
网络线程模型
首先需要理清broker的网络模型和异步处理框架,使用netty作为通信中间件,并在此基础上封装了一层qmq的网络通信模块,通过requestCode来关联processor,处理对应的逻辑。
线程模型基于netty的nio reactor模型。这部分体现在上图虚线框外面的部分。
注册sendMessage请求的processor

registerProcessor接收三个参数:
requestCode 与数据包中解析出来的对应
processor,用来处理对应数据包
处理线程池。
接收消息的处理器是SendMessageProcessor,不同的逻辑用不同的processor,很好的做到了解耦,而且代码很清晰干净。
启动NettyServer,设置bossGroup和workerGroup。

第80行添加通用serverHandler,通过requestCode路由processor。
serverHandler收到请求后交给第一步注册的processor在NettyRequestExecutor中来处理请求(这里从Netty的线程转到用户线程来执行,使用Netty时一定要注意,不要在netty的IO线程里执行IO等会block的操作)。

NettyRequestExecutor会调用对应的processor来处理业务逻辑,并且对processor返回值CompletableFuture添加callBack方法,callBack方法作用是回写对应的网络请求。

第92行添加callback,把msg接收的结果回写给producer。这一步被执行之后,producer端发起的sendMsg()将返回。对producer来说消息交付broker成功。
SendMessageProcessor将数据包反序列化为List<RawMessage>,然后交给SendMessageWorker来处理。

SendMessageWorker.receive如下:

第73行RawMessage转成ReceivingMessage,对每个Msg添加promise(其实就是将给每个Message绑定一个future,这个的作用就是这一批消息的future都完成了,这一批消息才算处理完成)。
第75行是处理单个消息的invoke方法,这一步会执行消息的各种校验,写入文件操作。
第78行Futures.transform利用guava类合并每个Message的结果(合并每个Message的future)。
以上解释了文章开头的1、2、3三个问题,下面分析消息的处理和问题4。
消息的处理流程
业务逻辑在invoke方法中,invoke使用的链式处理,可以进行扩展,最终主要逻辑在doInvoke方法中:

第104行校验Broker是否只读,比如Broker如果与metaserver之间失去心跳了,那么可能被标记为只读状态,只读状态不能接收消息。
第109行校验是否Broker的主从同步延迟是不是过大,超过一定的阈值Broker将被标记为只读。
第122行开始调用store层写入消息。
第123行是处理slave的同步(HA特性)。
注意:这个方法在返回之前,无论是校验不通过还是成功,最后都会调用RecivingMessage的done方法,完成自己的future。


第54行保证每个消息最后都会执行SettableFuture.set方法,并最终被第3步的第78行Futures.transform收集。
最终最后一个完成的消息会触发第4步骤的的callBack方法,回写结果。
HA模式下的高可用

messageStore.putMessage(message)之后,在第167行,如果消息可靠等级不是高级别的,直接返回,也就是不需要同步等待slave确认。
第179行,如果消息的级别是high,那需要将消息放入waitSlaveSyncQueue中(等待slave确认队列),等待slave的确认,只有在slave确认收到这个消息后(slave会通过另外一个主从同步端口从master拉取消息,每次拉取的时候都会带上上次拉取后的最大offset,那么小于等于这个offset的消息都被认为已经同步到了slave了),才会调用此消息的end操作,这个时候才会返回给producer,如下所示:

由此可见,高可用体现在两方面:
messageStore.putMessage(message)成功,即写入page cache成功
slave确认收到该消息
总结
以上就是broker接收消息的整个过程,良好的线程模型和异步框架是高性能的基础。HA机制也使得消息是可靠热备的。
但是上述过程只是保证了消息写入messageLog成功,如何刷盘,如何构建索引文件,以及如何通知订阅消息的consumer,还没有分析,留在以后再分析吧。




