十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

rocketMQ proxy架构分析

rocketMQ proxy架构分析 介绍proxy 模块是 RocketMQ 5.x 的 无状态代理组件 核心思路是对外提供 gRPC 协议面向多语言客户端对内把请求翻译成 Remoting 协议访问 Broker/Namesrv并原生支持 pop 消费模型。它有两种部署模式见 README.md Cluster 模式 Proxy 作为独立集群通过 RPC 与 Broker 通信存算分离。Local 模式 Proxy 与 Broker 同进程部署。两种模式通过 ServiceManagerFactory 切换 Local / Cluster 两套实现。各模块作用一、remoting / activityRemoting 协议入口的活动处理器remoting 包负责实现 自定义 Remoting 协议 的接入区别于 gRPC入口是MultiProtocolRemotingServer让老版 RocketMQ 4.x 客户端也能连到 Proxy。activity 子包里的类都实现 NettyRequestProcessor 是 按请求码路由的具体业务处理单元 类作用AbstractRemotingActivity基类。统一构造 ProxyContext 、执行 RequestPipeline 鉴权等、统一异常→响应码映射、回写响应ClientManagerActivity处理心跳/注册/注销把 producer/consumer 的 channel 注册进 MessagingProcessor 维护 RemotingChannelManagerConsumerManagerActivity消费者管理类查消费者列表/连接、锁/解锁 MQ、offset 查询更新等SendMessageActivity发送消息含 batch/消费端回退并做 topic 消息类型校验、事务消息订阅注册PopMessageActivitypop 拉取消息超时时间基于 pollTime 计算PullMessageActivity传统 pull 拉取AckMessageActivity / ChangeInvisibleTimeActivityACK 确认 / 修改消息不可见时间延迟重投GetTopicRouteActivity获取 topic 路由TransactionActivity结束事务commit/rollback其它子包pipeline RequestPipeline 责任链用于鉴权等前置处理。protocol 协议协商、SSL/TLS 多协议握手。channel RemotingChannel 及其管理。common RemotingConverter 协议转换。二、service / relay中转/透传服务ProxyRelayService 负责把 运维/管理类请求 消费进度查询、消费详情、事务状态回查转发到 Broker或从 Broker 侧直接回写给客户端。AbstractProxyRelayService 实现 processCheckTransactionState 回查事务状态时先落一条 TransactionData 再转发。LocalProxyRelayService Local 模式直接用 BrokerController 的 RemotingServer 把结果写回客户端如 getConsumerRunningInfo 、 consumeMessageDirectly ClusterProxyRelayService Cluster 模式 尚未实现 not implement yet 。辅助类 ProxyChannel 、 ProxyRelayResult 、 RelayData 用于承载中转结果。三、service / transaction事务消息支持TransactionService 负责事务消息的 订阅登记、事务数据管理、结束事务请求构造 。AbstractTransactionService 维护 TransactionDataManager 实现事务数据的新增/取出、生成 EndTransactionRequestHeader 。ClusterTransactionService Cluster 模式核心。通过 心跳 把事务生产者的 group 注册到对应 Broker否则 Broker 不认这个事务 group并维护 brokerAddr → brokerName 映射。内含 TxHeartbeatServiceThread 定时扫描发送心跳。LocalTransactionService Local 模式 空实现 因为 producer channel 已直接进 Broker 的 producerManager 无需 Proxy 额外处理。数据类 TransactionData 、 TransactionDataManager 、 EndTransactionRequestData 。四、service 其他子包子类作用message消息收发底层操作send/pop/pull/ack/changeInvisibleTime/offset/lock 等封装成对 Broker 的 RPCroute路由服务Caffeine 缓存路由、 MessageQueueSelector 读写队列选择含故障延迟 MQFaultStrategy metadata元数据topic 消息类型、订阅组配置receiptpop 消费的 ReceiptHandle 管理channelSimpleChannel / InvocationChannel 管理Local 模式进程内通道抽象clientProxy 侧的客户端 ClusterConsumerManager 、 ProxyClientRemotingProcessoradmin运维管理创建/更新 topic、订阅组等sysmessage系统消息同步Cluster 多 Proxy 间通过系统 topic 广播消费者心跳保证任一 Proxy 都能看到全量在线消费者ServiceManager 是这些服务的聚合门面统一暴露 MessageService / TopicRouteService / TransactionService / ProxyRelayService / MetadataService / AdminService 等。五、processor核心编排层MessagingProcessor 是 统一的业务门面 gRPC 和 Remoting 两个入口最终都调用它。它负责消息发送/消费编排、队列选择、事务结束、客户端producer/consumer注册管理、ReceiptHandle 管理等内部再委托给 ServiceManager。其子包channel RemoteChannel 跨 Proxy 的远程通道抽象及序列化。validator topic 消息类型校验器。TransactionProcessor 、 ProducerProcessor 、 ConsumerProcessor 等是分领域的具体实现。六、grpcgRPC 协议入口GrpcServer 提供 gRPC 服务 v2 子包是基于 rocketmq-apis protobuf的 MessagingService 实现 interceptor 做鉴权/上下文/异常处理。 AbstractMessingActivity 是 gRPC 侧各 Activity 的基类含 topic/group 校验。七、common / config / metrics基础支撑common ProxyContext 、 ReceiptHandleGroup 、 ProxyException 等公共模型。config ConfigurationManager 、 ProxyConfig 配置加载。metrics ProxyMetricsManager 监控指标。一句话总结 remoting/activity 是 Remoting 协议的业务入口 processor 是统一编排门面 service 是真正干活的后端消息/路由/事务/中转/元数据等三者构成「协议接入 → 编排 → 后端服务」三层结构 relay 管管理类请求的中转 transaction 管事务消息在 Proxy 侧的心跳注册与事务数据维护。
返回列表