高性能队列 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 系统中的应用,核心逻辑包括:
-
初始化 Disruptor:为每个 nodeId
创建独立的 RingBuffer。 -
按照节点转发消息:将用户消息存入对应节点的 RingBuffer。 -
消息处理: onEvent
方法从 RingBuffer 读取消息,并执行转发。
通过 Disruptor,可以大幅降低锁竞争,提高 IM 系统的吞吐量,使其能够在高并发环境下稳定运行。
4. 最后
欢迎关注加瓦点灯,每天推送干货知识,一起进步!
本文由 mdnice 多平台发布
相关文章:
高性能队列 Disruptor 在 IM 系统中的实战
高性能队列 Disruptor 在 IM 系统中的实战 前三期我们介绍了Disruptor的典型使用场景和相关高性能原理,本期我介绍一下Disruptor在IM系统用的应用实战,IM系统即社交聊天系统,对实时性的要求非常高,非常符合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
题目来源:buuctf BUU XSS COURSE 1 一、打开靶机,整理信息 有吐槽和登陆两个尝试点,题目名称提示是XSS漏洞 XSS(Cross-Site Scripting)漏洞 1.定义:跨站脚本攻击,是一种常见的 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来打开文件流是一个常见的操作,用于创建一个输入文件流,并尝试打开名为 question_list的文件。if (!in.is_open()):检查文件是否成功打开。如果文件未…...
“上门按摩” 小程序开发项目:基于 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,肯定要先学…...

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

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

【C++高并发服务器WebServer】-2:exec函数簇、进程控制
本文目录 一、exec函数簇介绍二、exec函数簇 一、exec函数簇介绍 exec 函数族的作用是根据指定的文件名找到可执行文件,并用它来取代调用进程的内容,换句话说,就是在调用进程内部执行一个可执行文件。 exec函数族的函数执行成功后不会返回&…...

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

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

SpringBoot打包为JAR包或WAR 包,这两种打包方式在运行时端口将如何采用?又有什么不同?这篇文章将给你解惑
你们好,我是金金金。 前提 SpringBoot打包方式:不是jar就是war包 场景 写这篇文章也是我遇到一个很不理解的点,所以就去研究了一下 场景如下: 这是后端生产配置文件给项目设置的端口,然后我前端写的url是:我就很纳闷,前端写了域名以及后端服务上下文路径,咋没写端口呢,…...

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

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

JavaWeb 学习笔记 XML 和 Json 篇 | 020
今日推荐语 愿你遇见好天气,愿你的征途铺满了星星——圣埃克苏佩里 日期 学习内容 打卡编号2025年01月23日JavaWeb笔记 XML 和 Json 篇020 前言 哈喽,我是菜鸟阿康。 以下是我的学习笔记,既做打卡也做分享,希望对你也有所帮助…...
在Raspbian上,如何获取树莓派的CPU当前频率
本文不用汇编实现,因为我是要用在 Go 里的,Go 并不支持内联汇编,要用汇编比较麻烦。而且项目并不是很在意性能,所以直接用命令获取内核准备好的。 在Raspbian上,CPU 信息存放在/sys/devices/system/cpu/中,…...
网络打印机的搜索与连接(一)
介绍 网络打印机就是可以通过网络连接上的打印机,这类打印机分2种:自身具有互联网接入功能可以分配IP的打印机我们称为网络打印机、另外一种就是被某台电脑连接上去后通过共享的方式共享到网络里面的我们称为共享打印机。现在还有一种可以通过互联网连接…...

LangChain + llamaFactory + Qwen2-7b-VL 构建本地RAG问答系统
单纯仅靠LLM会产生误导性的 “幻觉”,训练数据会过时,处理特定知识时效率不高,缺乏专业领域的深度洞察,同时在推理能力上也有所欠缺。 正是在这样的背景下,检索增强生成技术(Retrieval-Augmented Generati…...
ubuntu搭建nfs服务centos挂载访问
在Ubuntu上设置NFS服务器 在Ubuntu上,你可以使用apt包管理器来安装NFS服务器。打开终端并运行: sudo apt update sudo apt install nfs-kernel-server创建共享目录 创建一个目录用于共享,例如/shared: sudo mkdir /shared sud…...
R语言AI模型部署方案:精准离线运行详解
R语言AI模型部署方案:精准离线运行详解 一、项目概述 本文将构建一个完整的R语言AI部署解决方案,实现鸢尾花分类模型的训练、保存、离线部署和预测功能。核心特点: 100%离线运行能力自包含环境依赖生产级错误处理跨平台兼容性模型版本管理# 文件结构说明 Iris_AI_Deployme…...

中南大学无人机智能体的全面评估!BEDI:用于评估无人机上具身智能体的综合性基准测试
作者:Mingning Guo, Mengwei Wu, Jiarun He, Shaoxian Li, Haifeng Li, Chao Tao单位:中南大学地球科学与信息物理学院论文标题:BEDI: A Comprehensive Benchmark for Evaluating Embodied Agents on UAVs论文链接:https://arxiv.…...

剑指offer20_链表中环的入口节点
链表中环的入口节点 给定一个链表,若其中包含环,则输出环的入口节点。 若其中不包含环,则输出null。 数据范围 节点 val 值取值范围 [ 1 , 1000 ] [1,1000] [1,1000]。 节点 val 值各不相同。 链表长度 [ 0 , 500 ] [0,500] [0,500]。 …...

如何在看板中有效管理突发紧急任务
在看板中有效管理突发紧急任务需要:设立专门的紧急任务通道、重新调整任务优先级、保持适度的WIP(Work-in-Progress)弹性、优化任务处理流程、提高团队应对突发情况的敏捷性。其中,设立专门的紧急任务通道尤为重要,这能…...
Spring Boot面试题精选汇总
🤟致敬读者 🟩感谢阅读🟦笑口常开🟪生日快乐⬛早点睡觉 📘博主相关 🟧博主信息🟨博客首页🟫专栏推荐🟥活动信息 文章目录 Spring Boot面试题精选汇总⚙️ **一、核心概…...

微服务商城-商品微服务
数据表 CREATE TABLE product (id bigint(20) UNSIGNED NOT NULL AUTO_INCREMENT COMMENT 商品id,cateid smallint(6) UNSIGNED NOT NULL DEFAULT 0 COMMENT 类别Id,name varchar(100) NOT NULL DEFAULT COMMENT 商品名称,subtitle varchar(200) NOT NULL DEFAULT COMMENT 商…...

EtherNet/IP转DeviceNet协议网关详解
一,设备主要功能 疆鸿智能JH-DVN-EIP本产品是自主研发的一款EtherNet/IP从站功能的通讯网关。该产品主要功能是连接DeviceNet总线和EtherNet/IP网络,本网关连接到EtherNet/IP总线中做为从站使用,连接到DeviceNet总线中做为从站使用。 在自动…...
CRMEB 框架中 PHP 上传扩展开发:涵盖本地上传及阿里云 OSS、腾讯云 COS、七牛云
目前已有本地上传、阿里云OSS上传、腾讯云COS上传、七牛云上传扩展 扩展入口文件 文件目录 crmeb\services\upload\Upload.php namespace crmeb\services\upload;use crmeb\basic\BaseManager; use think\facade\Config;/*** Class Upload* package crmeb\services\upload* …...
实现弹窗随键盘上移居中
实现弹窗随键盘上移的核心思路 在Android中,可以通过监听键盘的显示和隐藏事件,动态调整弹窗的位置。关键点在于获取键盘高度,并计算剩余屏幕空间以重新定位弹窗。 // 在Activity或Fragment中设置键盘监听 val rootView findViewById<V…...