介绍
proxy 模块是 RocketMQ 5.x 的 无状态代理组件 ,核心思路是:对外提供 gRPC 协议(面向多语言客户端),对内把请求翻译成 Remoting 协议访问 Broker/Namesrv,并原生支持 pop 消费模型。它有两种部署模式(见 README.md ):
- Cluster 模式 :Proxy 作为独立集群,通过 RPC 与 Broker 通信(存算分离)。
- Local 模式 :Proxy 与 Broker 同进程部署。
两种模式通过 ServiceManagerFactory 切换 Local / Cluster 两套实现。
各模块作用
一、remoting / activity(Remoting 协议入口的活动处理器)
remoting 包负责实现 自定义 Remoting 协议 的接入(区别于 gRPC),入口是MultiProtocolRemotingServer,让老版 RocketMQ 4.x 客户端也能连到 Proxy。
activity 子包里的类都实现 NettyRequestProcessor ,是 按请求码路由的具体业务处理单元 :
| 类 | 作用 |
| AbstractRemotingActivity | 基类。统一构造 ProxyContext 、执行 RequestPipeline (鉴权等)、统一异常→响应码映射、回写响应 |
| ClientManagerActivity | 处理心跳/注册/注销:把 producer/consumer 的 channel 注册进 MessagingProcessor ,维护 RemotingChannelManager |
| ConsumerManagerActivity | 消费者管理类:查消费者列表/连接、锁/解锁 MQ、offset 查询更新等 |
| SendMessageActivity | 发送消息(含 batch/消费端回退),并做 topic 消息类型校验、事务消息订阅注册 |
| PopMessageActivity | pop 拉取消息,超时时间基于 pollTime 计算 |
| PullMessageActivity | 传统 pull 拉取 |
| AckMessageActivity / ChangeInvisibleTimeActivity | ACK 确认 / 修改消息不可见时间(延迟重投) |
| 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 的 RPC |
| route | 路由服务:Caffeine 缓存路由、 MessageQueueSelector (读写队列选择,含故障延迟 MQFaultStrategy ) |
| metadata | 元数据:topic 消息类型、订阅组配置 |
| receipt | pop 消费的 ReceiptHandle 管理 |
| channel | SimpleChannel / InvocationChannel 管理(Local 模式进程内通道抽象) |
| client | Proxy 侧的客户端: ClusterConsumerManager 、 ProxyClientRemotingProcessor |
| admin | 运维管理:创建/更新 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 等是分领域的具体实现。
六、grpc(gRPC 协议入口)
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 侧的心跳注册与事务数据维护。