一、什么时候需要救?ESB消息积压的应用场景
很多企业里都有这么一个工具,叫ESB,你可以把它理解成公司内部的“消息快递中转站”——各个系统之间要传的消息,比如电商里用户下单后,订单系统要告诉支付系统扣钱、要告诉库存系统减数量、要告诉营销系统给用户加积分,这些消息全都要经过ESB转递。平时业务量小的时候,ESB处理这些消息轻轻松松,就像平时快递点每天几十件件,很快就发完了。但一到大促,比如618、双11,或者突发的热点事件,用户下单量突然翻几十上百倍,ESB就会被消息“堵死”,就像快递点突然来了几万件件,堆成了小山,后面的消息没人管,有的订单支付了但没扣库存,有的用户下单了没收到通知,这就是消息积压,必须马上救!
二、核心急救方案:限流降级 + 优先级队列配合
救ESB的消息积压,不能硬扛,得用两个办法一起上,一个是“挡”,一个是“排”:
2.1 先“挡”:限流降级,给消息做“优先级筛选”
限流就像快递点的老板,高峰期说“今天最多收1000件件,超过的明天再送”,避免把整个快递点撑爆;降级就是说,有些没用的消息,比如发广告的营销短信,没人看的,高峰期就直接不发了,先把核心的活干了。
2.2 再“排”:优先级队列,给待处理的消息“排座次”
就算有消息要处理,也不能随便乱排,得把重要的先放前面,比如订单支付的消息肯定比发广告的重要,把这些消息放进优先级队列,核心的排前面,次核心的中间,非核心的后面,系统优先处理核心的,这样就不会因为非核心消息耽误核心业务。
三、实打实的落地示例:用Java快速实现(单一技术栈)
为了让大家能直接看懂甚至用到,我用Java写了一个完整的示例,完全基于JDK自带的工具,不用加额外的依赖,适合快速上线急救。首先明确技术栈:
// 技术栈:Java 17,使用JDK自带并发包的Semaphore实现限流,PriorityQueue实现优先级队列
然后是具体的代码,每一步都加了详细注释:
import java.util.PriorityQueue;
import java.util.concurrent.Semaphore;
// 1. 定义ESB消息类:每个消息带优先级,数值越小越重要(核心=1,次核心=2,非核心=3)
class EsbMessage implements Comparable<EsbMessage> {
private String messageId; // 每个消息的唯一ID,方便追踪
private int priority; // 优先级等级
private String content; // 消息内容,模拟实际业务数据
// 构造方法:创建消息时必须指定优先级
public EsbMessage(String messageId, int priority, String content) {
this.messageId = messageId;
this.priority = priority;
this.content = content;
}
// 2. 实现排序接口:让PriorityQueue按优先级从小到大排序(数值小的先出队)
@Override
public int compareTo(EsbMessage o) {
return Integer.compare(this.priority, o.priority);
}
// 提供get方法,方便代码中读取消息属性
public String getMessageId() { return messageId; }
public int getPriority() { return priority; }
public String getContent() { return content; }
}
// 3. ESB核心处理类:整合限流和优先级队列
public class EsbTrafficHandler {
// 限流核心:Semaphore是JDK的令牌工具,这里设为500,模拟ESB每秒最多处理500条消息
private static final Semaphore SYSTEM_LIMITER = new Semaphore(500);
// 优先级队列:存储所有待处理的消息,自动按优先级排序
private static final PriorityQueue<EsbMessage> PRIORITY_PROCESS_QUEUE = new PriorityQueue<>();
// 4. 入口方法:模拟ESB接收到新消息,先做限流处理
public static void onReceive(EsbMessage message) {
try {
// 尝试拿限流令牌:如果拿到,说明消息可以被处理;拿不到,进入降级逻辑
if (SYSTEM_LIMITER.tryAcquire()) {
// 拿到令牌,把消息放进优先级队列,等待处理
PRIORITY_PROCESS_QUEUE.offer(message);
System.out.printf("【正常处理】消息ID:%s,优先级:%d,已进入队列等待处理%n",
message.getMessageId(), message.getPriority());
} else {
// 拿不到令牌:触发降级逻辑,区分核心和非核心消息
if (message.getPriority() == 1) {
// 核心消息:不能丢,这里简化为存到本地临时日志(实际生产中用本地磁盘队列)
System.out.printf("【降级处理】核心消息ID:%s,临时存储到备用日志,后续重试%n",
message.getMessageId());
} else {
// 非核心消息:直接丢弃,属于降级,保障核心业务
System.out.printf("【降级丢弃】非核心消息ID:%s,因限流被丢弃%n",
message.getMessageId());
}
}
} catch (Exception e) {
e.printStackTrace();
}
}
// 5. 后台处理线程:从优先级队列取消息处理,模拟系统实际运行
public static void startProcessThread() {
new Thread(() -> {
while (true) { // 线程一直运行,随时处理队列里的消息
EsbMessage currentMessage = PRIORITY_PROCESS_QUEUE.poll(); // 取出队列头部的消息(优先级最高的)
if (currentMessage != null) {
try {
// 模拟处理消息的耗时:这里设100毫秒,实际中根据业务调整
Thread.sleep(100);
System.out.printf("【处理完成】消息ID:%s,优先级:%d,内容:%s%n",
currentMessage.getMessageId(), currentMessage.getPriority(), currentMessage.getContent());
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
// 处理完消息,释放限流令牌,让其他消息可以继续处理
SYSTEM_LIMITER.release();
}
}
}
}).start();
}
// 6. 测试主方法:模拟高峰期收到消息的场景
public static void main(String[] args) {
// 先启动后台处理线程,再模拟接收消息
startProcessThread();
// 模拟大促时收到10条消息:2条核心(订单/库存)、5条次核心(用户行为)、3条非核心(营销)
onReceive(new EsbMessage("B001", 1, "订单支付成功,更新订单状态"));
onReceive(new EsbMessage("B002", 1, "商品库存扣减成功,更新库存"));
onReceive(new EsbMessage("B003", 2, "用户下单,记录用户浏览行为"));
onReceive(new EsbMessage("B004", 2, "生成物流单号,通知快递"));
onReceive(new EsbMessage("B005", 2, "统计今日订单量,生成报表"));
onReceive(new EsbMessage("B006", 2, "奖励用户积分,触发营销活动"));
onReceive(new EsbMessage("B007", 2, "计算用户会员等级,生成权益"));
onReceive(new EsbMessage("B008", 3, "推送商品促销短信"));
onReceive(new EsbMessage("B009", 3, "APP推送新品到货通知"));
onReceive(new EsbMessage("B010", 3, "邮件发送节日活动邀请"));
}
}
这个代码跑起来后,你会发现,优先级1的核心消息会先被处理,非核心的直接被丢弃,系统不会崩,保证了核心业务的正常运行。
四、这个方案的优缺点,你要搞清楚
4.1 优点:救急快,效果明显
第一,能快速给系统“降压”,高峰期不会直接把ESB搞挂,避免更大的故障;第二,核心业务消息不会被淹没,比如订单支付的消息,就算限流了,也会优先处理,不会丢;第三,实现起来很简单,用JDK自带的工具就行,急救的时候十几分钟就能上线,不用搞复杂的架构。
4.2 缺点:不是完美的,要注意
第一,限流的阈值难调:设小了会丢太多核心消息,设大了还是会积压,平时要多测几个值;第二,非核心消息如果直接丢,会不会影响后续的营销效果?比如广告短信丢了,可能少点营收,这个要根据业务权衡;第三,核心消息临时存储如果用内存,万一内存不够,还是会丢,实际生产中要换成持久化的存储,比如本地文件或者Redis。
五、实战落地的注意事项,别踩坑
首先,限流阈值要按实际来:平时低峰期测一下ESB每秒能稳定处理多少条消息,高峰期把阈值设成平时的80%,留20%的余量应对突发;然后,优先级划分要准确:核心业务必须是直接影响交易的,比如订单、支付、库存,次核心是辅助体验的,非核心是对外营销的,别搞反了;还有,核心消息不能只存内存:刚才的示例里核心消息临时存在内存,实际中要存到本地磁盘或者分布式队列,比如用Redis的list存,防止内存溢出;另外,处理线程的数量要匹配CPU:不要开太多处理线程,不然线程切换会占资源,一般设成CPU核心数的2-4倍就行。
六、方案总结
ESB消息积压是流量高峰期常见的故障,这个急救方案的核心就是“先挡后排”:用限流挡住多余的消息,避免系统崩溃;用优先级队列保证核心业务消息优先处理,非核心消息要么降级丢弃要么后续重试,快速把系统从爆仓的状态拉回来,保住核心业务的稳定,不会影响公司的关键利益。这个方案适合快速急救,平时也可以作为ESB的标配配置,应对突发的流量峰值。
Comments