商会系统里有几件事做起来别扭:会费到期前要提醒会员,批量通知一次发几百条,报表导出要走几分钟,会员数据还要往另一个系统同步。它们共同的特点是耗时长、用户不需要当场拿到结果,但早期都写在接口里同步跑。
结果是导出报表时接口卡住,通知发一半失败要整批重来,提醒任务靠定时任务扫全表。后来把这些摘出来走消息队列,情况才好转。这篇记录用 RocketMQ 做的改造。
一、哪些该走消息
判断标准简单:用户不需要当场拿到结果的,都可以异步。我们确认要走的有三类。
第一类是耗时长的,比如报表导出、批量导入。接口只负责把任务丢进队列并返回任务号,用户去看进度。
第二类是需要重试的,比如短信和模板消息推送。第三方接口偶尔超时,同步做意味着用户跟着等,异步做可以在消费者里重试。
第三类是需要延时的,比如会费到期前七天提醒。用延时消息比定时扫表准,也不用维护一张待提醒表。
const {
Producer, Message } = require('ali-ons');
const producer = new Producer(config);
await producer.start();
async function submitExportTask(params) {
const taskId = genId();
await db.exportTask.create({
taskId, ...params, status: 'PENDING' });
await producer.send(new Message(TOPIC_EXPORT, 'export', JSON.stringify({
taskId, ...params })));
return taskId; // 接口立刻返回,前端拿 taskId 查进度
}
// 会费到期提醒:用延时消息提前七天投递
async function scheduleFeeRemind(feeId, dueAt) {
const delayMs = dueAt - Date.now() - 7 * 86400 * 1000;
if (delayMs <= 0) return;
const msg = new Message(TOPIC_REMIND, 'fee', JSON.stringify({
feeId }));
msg.setStartDeliverTime(Date.now() + delayMs);
await producer.send(msg);
}
二、消费者:幂等和失败处理
消费端容易出问题的地方是重复投递。网络抖动、消费者重启都会造成同一条消息被消费多次,所以消费逻辑必须幂等。
consumer.subscribe(TOPIC_EXPORT, 'export', async (msg) => {
const {
taskId, ...params } = JSON.parse(msg.body);
// 幂等:已处理过直接确认
const t = await db.exportTask.get(taskId);
if (!t || t.status !== 'PENDING') return 'CONSUME_SUCCESS';
try {
await db.exportTask.mark(taskId, 'RUNNING');
const fileUrl = await buildReport(params);
await db.exportTask.finish(taskId, fileUrl);
return 'CONSUME_SUCCESS';
} catch (e) {
await db.exportTask.mark(taskId, 'FAILED');
return 'RECONSUME_LATER'; // 交给队列重试
}
});
关键是状态机:任务有 PENDING / RUNNING / FINISHED / FAILED 几个状态,消费前先看状态,不是 PENDING 就直接确认,避免重复执行。失败返回重投,让队列按退避策略重试,次数用完进死信队列。
死信队列一定要配,不然失败的消息会一直重投,堵死消费者。
三、削峰
商会系统的流量极不均匀:年会报名开放那一分钟、会长在群里发了通知之后,请求会瞬间上来,平时又很闲。RocketMQ 在这里当缓冲区——请求先进队列,消费者按能承受的速度处理,超出部分排队,不会把数据库冲垮。
几个配置点:消费者线程数按下游承载能力定,不是越大越好;批量消费对通知类任务有效,一次取一批一起发能减少下游压力;队列数量决定并发上限,按峰值估算。
我们给不同类型分 Topic:导出类一个、通知类一个、同步类一个。这样互相不会被对方的堆积影响,导出任务跑得慢不会拖住通知。
四、踩过的坑
一是消息体不要塞大对象。有人把整个会员列表塞进消息体,结果超过大小限制。消息里只放标识和必要参数,数据让消费者去查。
二是别用消息替代事务。强一致场景还是要用事务消息或者本地消息表,普通消息保证不了。RocketMQ 的事务消息能用,但复杂度会上升,简单的做法反而是本地消息表加定时补偿。
三是消费进度监控要提前做。堆积条数和消费延迟没有监控,等用户反馈就晚了。
四是顺序。顺序消息性能差一些,多数场景不需要。会员状态变更按会员 ID 分区保证局部顺序就够了。
改造之后,导出接口从同步等待变成秒回,通知失败能自动重试,提醒任务也不用再扫全表。
我们这边在跑的商会管理系统叫未来漫城·商会互联平台,会费提醒、批量通知和报表导出都走 RocketMQ,年会报名那类峰值场景靠队列缓冲,数据库不再被打满。