模块:native-cloud-common-ringbuffer(包 com.huang.buffer

一、模块全景:一个模块,两条链路

ringbuffer 模块名字容易让人以为"只有环形队列",实际上它由两块内容组成:

  1. Lemon-Protocol 二进制 TCP 通信链路(纯 Netty):服务端收报文、反序列化、按业务标志分发;
  2. LemonBuffer 环形队列(LMAX Disruptor 4.0.0 封装):为通信链路触发的业务事件提供"异步化 + 削峰缓冲"的消费通道。

两者通过"事件"衔接:Netty 收到报文 → 分发到业务处理类 → 业务处理类把事件 publish 进环形队列 → Disruptor 消费线程异步处理。这就是"事件驱动"在通信模块中的落地形态。

flowchart LR subgraph "发送方" CW[ClientWorker
二进制发送客户端] end subgraph "通信链路 native-cloud-common-ringbuffer" NS[SocketServer
Netty 服务端] DEC[DelimiterBasedFrameDecoder
按换行符拆包] BYTE["ByteArrayDecoder
ByteBuf → byte[]"] GEN["GenerateMessageDataHandler
byte[] 转消息元数据列表"] META[MessageMetaDataHandler
按 businessFlag 分发] end subgraph "事件队列" LB[LemonBuffer
Disruptor 4.0.0] BH[BusinessHandler
业务处理类] EV[EventHandler
异步消费线程] end CW -- "Lemon-Protocol 二进制" --> NS NS --> DEC --> BYTE --> GEN --> META META -- "分发 doHandle" --> BH BH -- "publishEvent" --> LB LB -- "onEvent" --> EV

二、通信链路工作机制(Lemon-Protocol)

2.1 服务端 pipeline

SocketServer@Component,由 NettyStartListenerApplicationRunner)在应用启动后触发 init() 绑定端口。默认参数(ServerProperty.checkProperty() 兜底):

  • 端口 9798(可通过 lemon.ringbuffer.port 覆盖)
  • boss 线程 1、worker 线程 16

每个连接上的 pipeline:

DelimiterBasedFrameDecoder(65535, "\n")   ← 按换行符拆包,天然处理粘包/半包
        ↓
ByteArrayDecoder / ByteArrayEncoder       ← ByteBuf 与 byte[] 互转
        ↓
GenerateMessageDataHandler                ← 字节反序列化为 List<List<MessageMetaData>>
        ↓
MessageMetaDataHandler                    ← 按 businessFlag 分发到业务处理器

2.2 报文格式与反序列化

GenerateMessageDataHandlerByteTypeUtil.handleLineBytes() 完成反序列化,自定义报文格式为:

  • 1 字节:字段个数 n
  • n 字节:每个字段的类型(CHAR/SHORT/INTEGER/LONG/STRING/…)
  • 每字段 2 字节:字段值所占字节数
  • 数据区:各字段值
  • 尾部:参数名列表,用 | 分隔

解析结果是一条 List<MessageMetaData>,其中 MessageMetaData(type, data, name) 三元组,一次报文可以携带多条字段(甚至多组消息),减少小包数量。

2.3 事件分发:businessFlag → BusinessHandler

MessageMetaDataHandler 是链路的事件分发器:

  • 消息里带有 businessFlag 字段时,用它作为 key 查 handlers 映射;
  • handlersafterPropertiesSet() 里通过 applicationContext.getBeansOfType(BusinessHandler.class) 自动从 Spring 容器收集所有 BusinessHandler 实现,按 businessFlag() 建索引;
  • 命中后调用 businessHandler.doHandle(parameters),其中 parametersMap<String,Object>(字段名 → 值)。

业务扩展只需三步:实现 BusinessHandler 接口(businessFlag() + doHandle() + doRevert())→ 注册为 Spring bean → 在 BusinessFlag 枚举里登记标识。核心链路零改动。当前已登记:ENTRUST_BACK_WRITE(2200, 委托回写)FINANCING_PURCHASE(2201, 融资买入)

2.4 客户端与内部端口接口

  • ClientWorker:二进制协议发送客户端,start() 建连,publishData() 通过 channel.eventLoop().execute() 提交写任务(保证线程安全);
  • RingBufferInnerController:暴露 GET /ringBufferInner/port,供参数系统查询"业务子系统的二进制交互端口",用于动态建连。

三、LemonBuffer:事件驱动的核心

LemonBuffer<T extends MessageEvent> 是对 Disruptor 4.0.0 的薄封装,核心 API:

public class LemonBuffer<T extends MessageEvent> {
    private Disruptor<T> disruptor;

    // 构造:指定容量 + 事件工厂(DaemonThreadFactory 守护消费线程)
    public LemonBuffer(Integer bufferSize, EventFactory eventFactory);

    public void start();                                    // 启动消费线程
    public void shutdown();
    public void publishEvent(EventTranslator<T> translator);// 生产:翻译器在槽位内填充事件
    public EventHandlerGroup<T> handlerEvent(EventHandler<T>... handlers); // 消费:支持链式
    public Long nextSequence();                             // 低级 API,预留
}

事件对象预分配(零 GC 的关键)MessageEventFactory 抽象类实现 lmax 的 EventFactorynewInstance() 转调子类 generateMessageEvent()new Disruptor<>(eventFactory, bufferSize, ...) 在启动时会一次性预创建 1024 个事件对象填满整个环;此后每次发布只是"翻译器覆写已有槽位",不再产生新对象。

生产/消费模型publishEvent(EventTranslator) 由生产者(这里是 Netty IO 线程)执行,翻译器 translateTo(event, sequence) 在已分配槽位内写入业务字段;handlerEvent(EventHandler) 挂载的 onEvent(event, sequence, endOfBatch) 由 Disruptor 后台消费线程执行,endOfBatch 还可用于批量处理优化。

四、实际使用场景

4.1 场景一:委托回写异步化(credit-trade,唯一生产使用)

信用委托下单后,外部系统通过 Lemon-Protocol 二进制报文发来委托回写(businessFlag=2200),业务要求把委托状态从"待报/未报"更新为"已报",涉及 DB 查询 + 更新。若在 Netty worker 线程上同步执行 DB 操作,IO 线程吞吐会被数据库 RT 直接拖垮。

CreditEntrustBackWriteHandler 的做法——入队即返回

// 初始化:1024 槽位 + 预分配事件对象 + 启动守护消费线程
this.buffer = new LemonBuffer<>(1024, new EntrustWriteBackFactory());
this.buffer.handlerEvent(new EntrustWriteBackHandler());
this.buffer.start();

// Netty worker 线程执行:只翻译事件、发布、立即返回
public Result doHandle(Map<String, Object> parameters) {
    Long orderId = (Long) parameters.get(CommonConstant.ORDER_ID_KEY);
    buffer.publishEvent(new EventTranslator<EntrustWriteBackMassage>() {
        @Override
        public void translateTo(EntrustWriteBackMassage event, long sequence) {
            event.setOrderId(orderId);
        }
    });
    return new BizResult(ResultCodeEnum.WRITE_BACK_SUCCESS.getCode(), ...); // 立即返回
}

// Disruptor 守护消费线程执行:真正的 DB 逻辑
class EntrustWriteBackHandler implements EventHandler<EntrustWriteBackMassage> {
    public void onEvent(EntrustWriteBackMassage event, long sequence, boolean endOfBatch) {
        CreditEntrust creditEntrust = creditEntrustDomain.getCreditEntrust(event.getOrderId());
        // 撤单场景:先更新原委托状态
        // 待报/未报 → 已报
        creditEntrustDomain.updateCreditEntrust(creditEntrust);
    }
}
sequenceDiagram participant EXT as 外部系统 participant NETTY as Netty worker 线程 participant DISP as MessageMetaDataHandler participant DO as CreditEntrustBackWriteHandler.doHandle participant RB as LemonBuffer
(1024 槽位) participant TH as Disruptor 消费线程
EntrustWriteBackHandler participant DB as CreditEntrustDomain/DB EXT->>NETTY: 二进制报文(businessFlag=2200) NETTY->>NETTY: 拆包/反序列化 NETTY->>DISP: 按 businessFlag 分发 DISP->>DO: doHandle(parameters) DO->>RB: publishEvent(translateTo 填充 orderId) RB-->>DO: 立即返回 DO-->>NETTY: WRITE_BACK_SUCCESS(入队即返回) Note over RB,TH: Disruptor 后台消费 TH->>DB: getCreditEntrust(orderId) TH->>TH: 撤单处理 / 状态流转(待报→已报) TH->>DB: updateCreditEntrust

链路效果:Netty worker 线程从"同步等 DB"变成"入队即返回",DB 慢操作被隔离到独立消费线程,IO 吞吐不再受数据库 RT 拖累;回写高峰时的突发委托被 1024 槽位环形缓冲吸收,消费端按自身节奏处理。

4.2 场景二:参数系统二进制同步客户端(argument)

DataSyncWorker 与动态路由联动,把二进制通信客户端也做成了"事件驱动":

  1. Nacos 实例变更事件 → DynamicRoutesCache 更新路由 → DataSyncWorker.initWorkers()
  2. 对每个 application 的分片主节点 URI,GET http://{host}:{port}/ringBufferInner/port 拿二进制端口;
  3. new ClientWorker(new ClientProperty(host, port)) 建连,维护进 workerMap(按 application 维度)。

客户端侧同样遵循事件模型:publishData() 通过 channel.eventLoop().execute() 提交写任务,由 Netty EventLoop 串行化写操作,天然线程安全,且连接销毁/重建随路由事件增量完成,备节点同步由主节点 redo 兜底。

五、优势与效果

维度 机制 效果
无锁并发 Disruptor 环形缓冲 + 序列号 CAS 推进,单生产单消费无锁顺序读写 无锁竞争、无上下文切换开销,高吞吐低延迟
零 GC EventFactory 启动预分配 1024 个事件对象,运行时覆写复用 高频小消息场景不产生对象分配,GC 压力趋零
异步解耦 Netty IO 线程 publishEvent 即返回,DB 逻辑在消费线程执行 IO 线程吞吐不受 DB RT 拖累,响应延迟稳定
削峰缓冲 环形缓冲天然吸收突发流量,消费端按节奏消费 委托回写高峰不丢消息、不阻塞收包
事件驱动扩展 BusinessFlag + BusinessHandler 接口化 + Spring 容器自动收集 新增业务零侵入核心链路,三件套即可接入
守护消费线程 DaemonThreadFactory.INSTANCE 消费线程不阻塞 JVM 优雅退出
缓存行友好 Disruptor 内部 padding 隔离热字段(防伪共享) 多核环境下避免 MESI 缓存行抖动

六、关键代码位置索引

组件 方法 位置
SocketServer init() / doInit()(pipeline 装配) native-cloud-common-ringbuffer/.../buffer/server/SocketServer.java
NettyStartListener run() .../buffer/listener/NettyStartListener.java
GenerateMessageDataHandler channelRead() 反序列化 .../buffer/handler/GenerateMessageDataHandler.java
MessageMetaDataHandler channelRead() 分发 / afterPropertiesSet() 收集 handler .../buffer/handler/MessageMetaDataHandler.java
BusinessHandler businessFlag() / doHandle() / doRevert() .../buffer/core/BusinessHandler.java
LemonBuffer publishEvent() / handlerEvent() / start() .../buffer/core/LemonBuffer.java
MessageEventFactory newInstance()generateMessageEvent() .../buffer/core/MessageEventFactory.java
CreditEntrustBackWriteHandler afterPropertiesSet() 初始化 buffer / doHandle() 发布 native-cloud-credit-trade/.../credit/handler/CreditEntrustBackWriteHandler.java
ClientWorker start() / publishData() / submitWriteTask() .../buffer/worker/ClientWorker.java
RingBufferInnerController GET /ringBufferInner/port .../buffer/controller/RingBufferInnerController.java
ServerProperty lemon.ringbuffer.* 配置 .../buffer/properties/ServerProperty.java

七、结语与演进方向

当前 ringbuffer 模块对 Disruptor 的利用是"高性能无锁队列 + 异步化",真实收益集中在零 GC、无锁、IO 与 DB 解耦三件事上,对委托回写这类高频小消息场景是合理且够用的。

如果要进一步物尽其用,两个低成本方向:

  1. handlerEvent 返回的 EventHandlerGroup 支持链式——handleEventsWith(...).then(...) 可编排"回写 → 对账 → 通知"消费者依赖链,把单消费扩展为流水线;
  2. ServerProperty.bufferSize 目前是死配置——容量写死在调用方 new LemonBuffer<>(1024, ...),可以接入配置中心按场景调整;
  3. 回写确有吞吐瓶颈时,可暴露 WorkHandler 多消费分片 + 可配置 WaitStrategy(当前用默认策略)。

相关阅读:动态路由三组件启动时序分析DataSyncWorker 建连时序的上游依赖)