当前位置: 首页 > news >正文

高性能队列 Disruptor 在 IM 系统中的实战

高性能队列 Disruptor 在 IM 系统中的实战

前三期我们介绍了Disruptor的典型使用场景和相关高性能原理,本期我介绍一下Disruptor在IM系统用的应用实战,IM系统即社交聊天系统,对实时性的要求非常高,非常符合Disruptor的使用场景。

本篇文章将结合实际代码,介绍如何在 IM 系统中使用 Disruptor 进行高效的消息转发。

1. Disruptor 在 IM 系统中的作用

在 IM 系统中,用户 A 发送消息给 B、C、D 时,需要根据 B、C、D 所在的服务器节点进行分组,并将消息转发到对应的节点上。为了确保高吞吐量和低延迟,我们使用 Disruptor 作为高性能队列。

2. 代码实现

2.1 初始化 Disruptor

当某个节点 nodeId 还没有对应的 RingBuffer 时,我们需要创建一个新的 Disruptor,并将其存入 ringBufferMap 中。

   private final Map<String, RingBuffer<ClusterPublishEvent>> ringBufferMap = new ConcurrentHashMap<>();

    public ClusterQueueService(Server server) {
        this.mServer = server;
    }

    public void publishMessage(String nodeId, String fromUser, String clientId, String topic, byte[] payload) {
        if (!ringBufferMap.containsKey(nodeId)) {
            long st = System.currentTimeMillis();
            synchronized (ringBufferMap){
                if(!ringBufferMap.containsKey(nodeId)) {
                    BlockingWaitStrategy strategy = new BlockingWaitStrategy();
                    Disruptor<ClusterPublishEvent> disruptor = new Disruptor<>(
                        new ClusterPublishEventFactory(), 1024 * 1024, DaemonThreadFactory.INSTANCE,
                        ProducerType.SINGLE, strategy);
                    disruptor.handleEventsWith(new ClusterPublishEventHandler(mServer, nodeId));
                    disruptor.setDefaultExceptionHandler(new IgnoreExceptionHandler());
                    disruptor.start();
                    ringBufferMap.put(nodeId, disruptor.getRingBuffer());
                }
            }
            log.info("publishMessage create RingBuffer cost:{}ms, ringBufferMap:{},size:{}", System.currentTimeMillis() - st, ringBufferMap, ringBufferMap.size());
        }
        RingBuffer<ClusterPublishEvent> ringBuffer = ringBufferMap.get(nodeId);
        long sequence = ringBuffer.next();
        // 当环形缓冲区未用完时, 返回的是空对象,否则,返回的是缓存的数据。
        ClusterPublishEvent clusterEvent = ringBuffer.get(sequence);
        clusterEvent.setFromUser(fromUser);
        clusterEvent.setClientId(clientId);
        // 此topic,是节点转发的topic: NM2R, NTF,DESTROYUSER, 只有这三种
        clusterEvent.setTopic(topic);
        clusterEvent.setPayload(payload);
        clusterEvent.setTraceId(MDC.get(ImSvcConstants.TRACE_ID));
        // 发布事件, 会触发ClusterPublishEventHandler.onEvent方法
        ringBuffer.publish(sequence);
    }

关键点解析:

  • 采用 BlockingWaitStrategy 作为等待策略,确保高效的 CPU 资源利用。
  • 采用 DaemonThreadFactory.INSTANCE 创建线程池,避免应用程序退出时线程未正常回收。
  • handleEventsWith 设定事件处理器 ClusterPublishEventHandler,用于消息处理。
  • setDefaultExceptionHandler 避免异常影响消息处理流程。

2.2 按照节点转发消息

根据用户所在的服务节点,进行消息转发(发送消息事件到Disruptor)

    public void publish2Receivers(Long messageId, Set<String> receivers, String exceptClientId, int pullType, String topic) {
        //未绑定broker的用户默认由本中心处理
        Map<String, String> allReceiverMap = new HashMap<>();
        for (String receiver : receivers) {
            allReceiverMap.put(receiver, localNodeId);
        }
        //从分布式缓存获取获取用户路由
        Map<String, String> receiverMap = userRouteStore.getAll(receivers);
        allReceiverMap.putAll(receiverMap);
        Map<String, Set<String>> nodeMap = new HashMap<>();
        //使用nodeId分组
        allReceiverMap.forEach((receiver, nodeId) -> {
            if (!nodeMap.containsKey(nodeId)) {
                nodeMap.put(nodeId, new HashSet<>());
            }
            nodeMap.get(nodeId).add(receiver);
        });
        //获取可用节点
        Cluster cluster = mServer.getHazelcastInstance().getCluster();
        Set<Member> members = cluster.getMembers();
        List<String> collect = members.stream().map(member -> member.getStringAttribute(HZ_Cluster_Node_ID)).collect(Collectors.toList());
        log.info("hazelcast node list:{}",JSON.toJSONString(collect));
        Map<String, Member> memberMap = members.stream().collect(Collectors.toMap(
            member -> member.getStringAttribute(HZ_Cluster_Node_ID), member -> member, (k1, k2 )->k1));
        //按照节点分发
        nodeMap.forEach((nodeId, set) -> {
            // 转发到其他节点发送
            if (!nodeId.equals(localNodeId) && memberMap.containsKey(nodeId)) {
                WFCMessage.NotifyMessage2Receivers notifyMessage2Receivers = WFCMessage.NotifyMessage2Receivers.newBuilder()
                    .setMessageId(messageId)
                    .addAllReceivers(set)
                    .setExceptClientId(exceptClientId==null?"":exceptClientId)
                    .setPullType(pullType)
                    .setTopic(topic)
                    .build();
                clusterQueueService.publishMessage(nodeId,nodeId,null, IMTopic.NotifyMessage2ReceiversTopic, notifyMessage2Receivers.toByteArray());
            }
            // 当前节点处理发送
            else {
                WFCMessage.Message message = mServer.getStore().messagesStore().getMessage(messageId);

                if (message != null) {
                    // Add By Youqibin 16:11 2022/3/15 接收通知前置处理
                    preHandle(message, set);
                    mServer.getImBusinessScheduler().execute(() ->messagesPublisher.publish2Receivers(message, set, exceptClientId, pullType, localNodeId));
                    // Add By Youqibin 16:11 2022/3/15 接收通知后置处理
                    postHandle(message, set);
                }
            }
        });
    }

关键点解析:

  • clusterQueueService.publishMessage, 使用Disruptor发送消息事件,高性能异步处理

2.3 事件处理器 onEvent

当 Disruptor 事件发布后,ClusterPublishEventHandler.onEvent 负责实际的消息转发逻辑。

public class ClusterPublishEventHandler implements EventHandler<ClusterPublishEvent{
    private final Server server;
    private final String nodeId;

    public ClusterPublishEventHandler(Server server, String nodeId) {
        this.server = server;
        this.nodeId = nodeId;
    }

    @Override
    public void onEvent(ClusterPublishEvent event, long sequence, boolean endOfBatch) {
        log.info("Processing event: {} on node: {}", event, nodeId);
        server.forwardMessage(nodeId, event.getFromUser(), event.getClientId(), event.getTopic(), event.getPayload());
    }
}

关键点解析:

  • onEvent 方法接收到 ClusterPublishEvent 后,调用 server.forwardMessage 进行消息转发。
  • endOfBatch 用于标识当前事件是否为批处理中的最后一个事件。
  • log.info 记录消息处理的关键日志,便于后续排查。

3. 总结

本文介绍了 Disruptor 在 IM 系统中的应用,核心逻辑包括:

  1. 初始化 Disruptor:为每个 nodeId 创建独立的 RingBuffer。
  2. 按照节点转发消息:将用户消息存入对应节点的 RingBuffer。
  3. 消息处理onEvent 方法从 RingBuffer 读取消息,并执行转发。

通过 Disruptor,可以大幅降低锁竞争,提高 IM 系统的吞吐量,使其能够在高并发环境下稳定运行。

4. 最后

欢迎关注加瓦点灯,每天推送干货知识,一起进步!

本文由 mdnice 多平台发布

相关文章:

高性能队列 Disruptor 在 IM 系统中的实战

高性能队列 Disruptor 在 IM 系统中的实战 前三期我们介绍了Disruptor的典型使用场景和相关高性能原理&#xff0c;本期我介绍一下Disruptor在IM系统用的应用实战&#xff0c;IM系统即社交聊天系统&#xff0c;对实时性的要求非常高&#xff0c;非常符合Disruptor的使用场景。 …...

原生HTML集合

一、表格 1、固定表格 <div class"tablebox"><div class"table-container"><table id"myTable" border"0" cellspacing"0" cellpadding"0"><thead><tr></tr></thead>…...

ES6 简单练习笔记--变量申明

一、ES5 变量定义 1.在全局作用域中 this 其实就是window对象 <script>console.log(window this) </script>输出结果: true 2.在全局作用域中用var定义一个变量其实就相当于在window上定义了一个属性 例如: var name "孙悟空" 其实就相当于执行了 win…...

2025.1.21——六、BUU XSS COURSE 1

题目来源&#xff1a;buuctf BUU XSS COURSE 1 一、打开靶机&#xff0c;整理信息 有吐槽和登陆两个尝试点&#xff0c;题目名称提示是XSS漏洞 XSS&#xff08;Cross-Site Scripting&#xff09;漏洞 1.定义&#xff1a;跨站脚本攻击&#xff0c;是一种常见的 Web 安全漏洞。攻…...

Linux - 五种常见I/O模型

I/O操作 (输入/输出操作, Input/Output) 是指计算机与外部设备就行数据交互的过程. 什么是外部设备: 如键盘, 鼠标, 硬盘, 网卡等. 五种常见的 I/O 模型: 阻塞 I/O非阻塞 I/O信号驱动 I/OI/O 多路复用异步 I/O 阻塞 I/O 阻塞 I/O 的特点: 当用户发起 I/O 请求后, 进程/线程就…...

【负载均衡式在线OJ】加载题目信息(文件版)

目录 如何读取文件 -- 常见流程 代码 如何读取文件 -- 常见流程 在C中使用 std::ifstream来打开文件流是一个常见的操作&#xff0c;用于创建一个输入文件流&#xff0c;并尝试打开名为 question_list的文件。if (!in.is_open())&#xff1a;检查文件是否成功打开。如果文件未…...

“上门按摩” 小程序开发项目:基于 SOP 的全流程管理

在竞争激烈的生活服务市场,“上门按摩” 服务需求日益增长。为满足这一需求,我们启动了 O2O 模式的 “上门按摩” 小程序开发项目,该项目涵盖客户端、系统管理端、技师端。以下将通过各类 SOP 对项目进行全面管理,确保项目顺利推进。 一、项目启动 SOP:5W2H 分析法 What(…...

WPF1-从最简单的xaml开始

1. 最简单的WPF应用 1.1. App.config1.2. App.xaml 和 App.xaml.cs1.3. MainWindow.xaml 和 MainWindow.xaml.cs 2. 正式开始分析 2.1. 声明即定义2.2. 命名空间 2.2.1. xaml的Property和Attribute2.2.2. xaml中命名空间2.2.3. partial关键字 学习WPF&#xff0c;肯定要先学…...

2025牛客寒假算法营2

A题 知识点&#xff1a;模拟 打卡。检查给定的七个整数是否仅包含 1,2,3,5,6 即可。为了便于书写&#xff0c;我们可以反过来&#xff0c;检查这七个整数是否不为 4 和 7。 时间 O(1)&#xff1b;空间 O(1)。 #include <bits/stdc.h> using namespace std;signed main()…...

编译Android平台使用的FFmpeg库

目录 前言 一、编译环境 二、搭建环境 1.安装MSYS2 2.更新系统包 2.1 打开MSYS2 MinGW 64-bit终端&#xff08;mingw64.exe&#xff09; 2.2 更新所有软件包到最新版本 2.3 安装必要的工具和库。 3. 克隆FFmpeg源码 4. 配置编译选项 5. 执行编译 总结 前言 记录学习…...

【C++高并发服务器WebServer】-2:exec函数簇、进程控制

本文目录 一、exec函数簇介绍二、exec函数簇 一、exec函数簇介绍 exec 函数族的作用是根据指定的文件名找到可执行文件&#xff0c;并用它来取代调用进程的内容&#xff0c;换句话说&#xff0c;就是在调用进程内部执行一个可执行文件。 exec函数族的函数执行成功后不会返回&…...

力扣707题(2)——设计链表

#题目 #3,5和6的代码 今天看剩下几个题的代码&#xff0c;1,2,4的代码已经在上篇博客写过了想看的小伙伴移步到&#xff1a; 力扣707题——设计链表-CSDN博客 //第3题头插法 void addAtHead(int val){ //记录头结点ListNode nhead; //新节点的创建,并让它指向原本头结点的后…...

K8S中ingress详解

Ingress介绍 Kubernetes 集群中&#xff0c;服务&#xff08;Service&#xff09;是一种抽象&#xff0c;它定义了一种访问 Pod 的方式&#xff0c;无论这些 Pod 如何变化&#xff0c;服务都保持不变。服务可以被映射到一个静态的 IP 地址&#xff08;ClusterIP&#xff09;、一…...

SpringBoot打包为JAR包或WAR 包,这两种打包方式在运行时端口将如何采用?又有什么不同?这篇文章将给你解惑

你们好,我是金金金。 前提 SpringBoot打包方式:不是jar就是war包 场景 写这篇文章也是我遇到一个很不理解的点,所以就去研究了一下 场景如下: 这是后端生产配置文件给项目设置的端口,然后我前端写的url是:我就很纳闷,前端写了域名以及后端服务上下文路径,咋没写端口呢,…...

zabbix6.0安装及常用监控配置

文章目录 部署zabbix-serverzabbix监控节点部署解决zabbix中文乱码创建主机组创建模版配置主机与模版关联 监控boot分区监控网卡流量出网卡流量监控进入和出的总流量监控内存监控服务器端口用户自定应监控key值 (监控mysql查询数量)zabbix触发器监控cpu监控入网卡流量 邮件告警…...

SQL-leetcode—1179. 重新格式化部门表

1179. 重新格式化部门表 表 Department&#xff1a; ---------------------- | Column Name | Type | ---------------------- | id | int | | revenue | int | | month | varchar | ---------------------- 在 SQL 中&#xff0c;(id, month) 是表的联合主键。 这个表格有关…...

JavaWeb 学习笔记 XML 和 Json 篇 | 020

今日推荐语 愿你遇见好天气,愿你的征途铺满了星星——圣埃克苏佩里 日期 学习内容 打卡编号2025年01月23日JavaWeb笔记 XML 和 Json 篇020 前言 哈喽&#xff0c;我是菜鸟阿康。 以下是我的学习笔记&#xff0c;既做打卡也做分享&#xff0c;希望对你也有所帮助…...

在Raspbian上,如何获取树莓派的CPU当前频率

本文不用汇编实现&#xff0c;因为我是要用在 Go 里的&#xff0c;Go 并不支持内联汇编&#xff0c;要用汇编比较麻烦。而且项目并不是很在意性能&#xff0c;所以直接用命令获取内核准备好的。 在Raspbian上&#xff0c;CPU 信息存放在/sys/devices/system/cpu/中&#xff0c…...

网络打印机的搜索与连接(一)

介绍 网络打印机就是可以通过网络连接上的打印机&#xff0c;这类打印机分2种&#xff1a;自身具有互联网接入功能可以分配IP的打印机我们称为网络打印机、另外一种就是被某台电脑连接上去后通过共享的方式共享到网络里面的我们称为共享打印机。现在还有一种可以通过互联网连接…...

LangChain + llamaFactory + Qwen2-7b-VL 构建本地RAG问答系统

单纯仅靠LLM会产生误导性的 “幻觉”&#xff0c;训练数据会过时&#xff0c;处理特定知识时效率不高&#xff0c;缺乏专业领域的深度洞察&#xff0c;同时在推理能力上也有所欠缺。 正是在这样的背景下&#xff0c;检索增强生成技术&#xff08;Retrieval-Augmented Generati…...

【自然语言处理(NLP)】介绍、发展史

文章目录 介绍发展史1. 规则驱动时期&#xff08;20世纪50年代-80年代&#xff09;技术特点标志性成果 2. 统计方法兴起&#xff08;1990年代-2000年代&#xff09;技术特点标志性成果 3. 神经网络复兴&#xff08;2010年代初至今&#xff09;技术特点标志性成果 4. 集成与应用…...

1.CSS的三大特性

css有三个非常重要的三个特性&#xff1a;层叠性、继承性、优先级 1.1 层叠性 想通选择器给设置想听的样式&#xff0c;此时一个样式就会覆盖&#xff08;层叠&#xff09;另一个冲突的样式。层叠性主要是解决样式冲突的问题。 <!DOCTYPE html> <html lang"en&…...

【分布式日志篇】从工具选型到实战部署:全面解析日志采集与管理路径

网罗开发 &#xff08;小红书、快手、视频号同名&#xff09; 大家好&#xff0c;我是 展菲&#xff0c;目前在上市企业从事人工智能项目研发管理工作&#xff0c;平时热衷于分享各种编程领域的软硬技能知识以及前沿技术&#xff0c;包括iOS、前端、Harmony OS、Java、Python等…...

基于springcloud汽车信息分析与可视化系统

基于Spring Cloud的汽车信息分析与可视化系统是一款旨在整合、分析汽车相关数据并以直观可视化方式呈现的应用系统。 一、系统架构 该系统基于先进的Spring Cloud架构构建&#xff0c;充分利用其分布式、微服务特性&#xff0c;确保系统具备高可用性、可扩展性和灵活性。Spri…...

TOGAF之架构标准规范-信息系统架构 | 数据架构

TOGAF是工业级的企业架构标准规范&#xff0c;信息系统架构阶段是由数据架构阶段以及应用架构阶段构成&#xff0c;本文主要描述信息系统架构阶段中的数据架构阶段。 如上所示&#xff0c;信息系统架构&#xff08;Information Systems Architectures&#xff09;在TOGAF标准规…...

Databend x 沉浸式翻译 | 基于 Databend Cloud 构建高效低成本的业务数据分析体系

「沉浸式翻译」是一个非常流行的双语对照网页翻译扩展工具&#xff0c;用户可以用它来即时翻译外文网页、PDF 文档、ePub 电子书、字幕等。它不仅可以实现原文加译文实时双语对照显示&#xff0c;还支持 Google、OpenAI、DeepL、微软、Gemini、Claude 等数十家翻译平台服务的自…...

cuda的并行运算介绍

cuda是如何使用GPU并行运算的&#xff1a; 以一个函数为例&#xff1a; duplicateWithKeys << <(P 255) / 256, 256 >> > (P,geomState.means2D,geomState.depths,geomState.point_offsets,binningState.point_list_keys_unsorted,binningState.point_list_…...

「全网最细 + 实战源码案例」设计模式——抽象工厂模式

核心思想 抽象工厂模式是一种创建型设计模式&#xff0c;它提供一个接口&#xff0c;用于创建一系列相关或互相依赖的对象&#xff0c;而无需指定它们的具体类。抽象工厂模式解决了产品族的问题&#xff0c;可以管理和创建一组相关的产品。 结构 1. 抽象工厂 定义创建一些列…...

领域驱动设计(DDD)四 订单管理系统实践步骤

以下是基于 领域驱动设计&#xff08;DDD&#xff09; 的订单管理系统实践步骤&#xff0c;系统功能主要包括订单的创建、更新、查询和状态管理&#xff0c;采用 Spring Boot 框架进行实现。 1. 需求分析 订单管理系统的基本功能&#xff1a; 订单创建&#xff1a;用户下单创…...

leetcode 面试经典 150 题:简化路径

链接简化路径题序号71题型字符串解法栈难度中等熟练度✅✅✅ 题目 给你一个字符串 path &#xff0c;表示指向某一文件或目录的 Unix 风格 绝对路径 &#xff08;以 ‘/’ 开头&#xff09;&#xff0c;请你将其转化为 更加简洁的规范路径。 在 Unix 风格的文件系统中规则如下…...