阿里云消息队列 RocketMQ:给商会系统做异步解耦与削峰的一次实践

简介: 本文介绍商会系统通过引入 RocketMQ 实现异步化改造:将耗时、可重试、需延时的任务(如报表导出、批量通知、会费提醒)从同步接口剥离,解耦执行与响应。通过消息队列实现削峰、幂等消费、失败重试与死信处理,显著提升系统稳定性与用户体验。(239字)

商会系统里有几件事做起来别扭:会费到期前要提醒会员,批量通知一次发几百条,报表导出要走几分钟,会员数据还要往另一个系统同步。它们共同的特点是耗时长、用户不需要当场拿到结果,但早期都写在接口里同步跑。

结果是导出报表时接口卡住,通知发一半失败要整批重来,提醒任务靠定时任务扫全表。后来把这些摘出来走消息队列,情况才好转。这篇记录用 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,年会报名那类峰值场景靠队列缓冲,数据库不再被打满。

相关文章
|
4月前
|
存储 人工智能 运维
十大 AI Agent Memory记忆系统全维度横评 主流方案架构、性能与场景选型指南
随着AI Agent从基础问答工具进化为可执行复杂长周期任务的智能体,记忆能力已经成为决定智能体上限的核心要素。传统基于向量数据库与RAG检索的技术方案,仅能实现简单信息检索,并不具备完整的记忆管理能力,在时序追踪、多代理一致性、分层存储、智能路由等方面存在明显短板。在实际生产环境中,大量AI Agent将近八成以上的计算资源消耗在重复梳理上下文信息上,真正用于业务执行的资源占比极低,这也是当前智能体规模化落地的核心瓶颈。
2136 2
idea实现protobuf的.proto文件编译成.java文件教程
1..proto文件语法高亮显示1.1 打开idea的插件列表1.2 下载protobuf辅助插件1.3 安装好后重启idea 在项目中新增配置生成环境 1.6.1
14834 0
|
11天前
|
人工智能 运维 安全
剧透丨AI 原生研发组织的探索和实践
9月23日下午,欢迎来到杭州国际博览中心一期一楼 103B 厅,与我们一起讨论 Agent 进入生产链路之后,研发组织将如何真正发生变化。
|
11天前
|
人工智能 PyTorch 云栖大会
亮点抢先看!AMD、中兴通讯等企业大咖解码下一代 OS|2026 云栖大会
一览来自 AMD、中兴通讯、阿里云、英特尔,以及上海交通大学等企业/高校专家的精彩金句和重磅议题。
|
17天前
|
人工智能 自然语言处理 文字识别
【2026.09.07~09.13】阿里云百炼最新产品动态周报
Token Plan个人版新增12类Agent Harness工具,支持搜索、图像生成等能力;控制台界面焕新,应用广场上新4大解决方案;MCP广场新增22个服务,覆盖金融、法律、OCR等场景;DeepSeek-V4.1-Flash模型上线,支持1M上下文与多模态理解。
326 0
|
7月前
|
人工智能 安全 API
AI数字员工落地:OpenClaw阿里云/本地企业级部署与千问/Coding Plan API配置指南
2026年AI智能体技术已进入规模化落地阶段,以Qwen3系列为代表的大模型具备超长文本上下文、多模态理解与多步骤复杂任务执行能力,让AI从辅助工具进化为可独立执行业务流程的数字员工。阿里云百炼平台推出的qwen3-max-2026-01-23模型,融合自主思考、联网检索、代码解析能力,在企业场景中实现高精度决策与任务处理。面对这一趋势,企业需要轻量化、安全可控、易扩展的AI智能体管理平台,OpenClaw(Clawdbot)正是支撑数字员工部署、调度、运维的核心框架,具备多平台接入、模块化技能、定时任务、安全沙箱四大核心能力,可帮助企业快速搭建AI自动化工作流,实现降本增效。本文完整讲解Op
1106 3
|
人工智能 JavaScript 前端开发
LangGraph架构解析
本文深入解析了传统Agent开发的三大痛点:状态管理碎片化、流程控制复杂及扩展性差,提出使用LangGraph通过有向图模型重构工作流,将LLM调用与工具执行抽象为节点,实现动态流程跳转。文中详述LangGraph四大核心组件——状态机引擎、节点设计、条件边与工具层集成,并结合生产环境最佳实践,如可视化调试、状态持久化与人工干预机制,最终对比LangGraph与传统方案的性能差异,给出选型建议。
2765 1
|
10月前
|
XML 算法 安全
详解RAG五种分块策略,技术原理、优劣对比与场景选型之道
本文详解RAG系统中五种核心分块策略——固定大小、语义、递归、基于文档结构及基于LLM的分块,涵盖其技术原理、优劣对比与适用场景。分块策略直接影响检索精度与生成质量,是构建高效RAG系统的关键。文章结合实例与决策树,指导读者根据文档类型与业务需求选择最优方案,并探讨当前挑战与前沿优化方向。
|
12月前
|
消息中间件 监控 Kubernetes
别再乱排查了!Kafka 消息积压、重复、丢失,根源基本都是 Rebalance!
大家好,我是小富~分享一次Kafka消息积压排查经历:消费者组因Rebalance导致消费能力骤降。本文详解Rebalance触发场景(消费者变更、分区扩容、订阅变化、超时等),剖析其引发的消息积压、重复消费、丢失等问题根源,并提供优化方案:调优超时参数、手动提交offset、启用粘性分配策略、保障消费幂等性。掌握这些,轻松应对Kafka常见故障!
2058 0

热门文章

最新文章