
后端微服务云原生【免费下载链接】midway A Node.js Serverless Framework for front-end/full-stack developers. Build the application for next decade. Works on AWS, Alibaba Cloud, Tencent Cloud and traditional VM/Container. Super easy integrate with React and Vue. 项目地址https://gitcode.com/gh_mirrors/mi/midway点击查看免费下载导读midwayjs/rabbitmq是 Midway 框架内置的 RabbitMQ 订阅/消费组件专门解决微服务通过消息队列解耦的场景——服务 A 负责向队列投递消息服务 B 负责消费队列中的任务。本文以该组件的 CHANGELOG.md 为时间线骨架完整梳理其从诞生到当前的演进脉络并结合仓库源码连接管理、消费者装载、链路追踪、本地 Mock 测试与官方文档带你从会用走向理解其内部实现。读完你将掌握如何独立部署或附加接入该组件、如何用ConsumerRabbitMQListener订阅队列、如何配置连接与各种 Exchange 路由策略、如何在无真实 MQ 环境下用midwayjs/mock完成本地测试以及它的断线重连、守卫与链路追踪等可靠性机制是怎么实现的。组件定位与可用性在开始编码之前先明确midwayjs/rabbitmq在 Midway 生态中的定位。根据 site/docs/extensions/rabbitmq.md 中的能力矩阵该组件能力支持情况可用于标准项目✅可用于 Serverless❌可用于一体化全栈应用✅包含独立主框架✅包含独立日志❌它既可以作为独立主框架直接启动一个纯消息消费进程也可以作为附加组件挂在midwayjs/koa、midwayjs/express等 Web 框架之下与 HTTP 服务共存于同一个应用内。从 package.json 可以看到当前仓库中该包的版本为4.2.3底层基于amqplib与amqp-connection-manager构建amqp-connection-manager4.1.15是运行时依赖amqplib与types/amqplib以 peerDependencies 形式要求使用方自行安装Node 引擎要求20。版本演进时间线从 2.7.0 到 4.x 的关键变更CHANGELOG 是理解组件能力边界与可靠性演进最权威的第一手资料。以下按时间倒序提炼出与 RabbitMQ 组件直接相关的关键条目版本时间关键变更含义3.7.02022-10-29依赖升级amqp-connection-manager至 v4.1.9保持连接管理库处于较新版本3.6.02022-10-10新增 guard 支持依赖升级至 v4.1.7消费方法可接入 Midway 守卫体系3.4.52022-07-25修复 rabbitmq 断线问题优化断线场景下的连接状态处理3.4.0-beta.112022-07-19连接失败时抛错启动阶段连接失败不再静默而是向上抛出异常3.0.132022-03-01修复 rabbitmq 配置 key统一配置读取键名3.0.0-beta.162022-01-11依赖升级amqp-connection-manager至 v4主版本升级带来新的连接管理 API2.11.02021-06-10新增 reconnection断线重连与对应测试组件具备自动重连能力2.7.02021-01-27新增 rabbitmq 组件组件首次引入 Midway 生态将 CHANGELOG 与源码对照可以提炼出三条核心演进主线连接管理现代化从 v4 版本amqp-connection-manager开始组件全面采用其连接池/自动重连模型并通过持续的小版本升级v4.1.0 → v4.1.9 → v4.1.15跟进上游修复。启动语义收紧早期版本连接失败时行为不明确3.4.0-beta.11起明确在启动阶段连接失败即抛错避免假启动。能力扩展2.11.0引入断线重连3.6.0引入 guard守卫使消费者与 Web 层保持一致的拦截能力。此外组件早在 v2 时代就以egg-rabbitmq-plus为蓝本改造而来——mq.ts 文件头部的注释明确标注了这段历史渊源。快速上手安装与开启组件安装依赖组件本身不捆绑amqplib需要显式安装$ npm i midwayjs/rabbitmq4 --save $ npm i amqplib --save $ npm i types/amqplib --save-dev或直接在 package.json 中声明{ dependencies: { midwayjs/rabbitmq: ^4.0.0, amqplib: ^0.10.1 }, devDependencies: { types/amqplib: ^0.8.2 } }作为独立主框架开启// src/configuration.ts import { Configuration } from midwayjs/core; import * as rabbitmq from midwayjs/rabbitmq; Configuration({ imports: [rabbitmq], }) export class MainConfiguration { async onReady() { // ... } }附加在 Web 主框架下// src/configuration.ts import { Configuration } from midwayjs/core; import * as koa from midwayjs/koa; import * as rabbitmq from midwayjs/rabbitmq; Configuration({ imports: [koa, rabbitmq], }) export class MainConfiguration { async onReady() { // ... } }开启组件后框架会读取配置键rabbitmq——这一点在 framework.ts 中体现得最直接configure() { return this.configService.getConfiguration(rabbitmq); }消费者Consumer实战目录结构约定官方推荐将消费者集中在src/consumer目录下my_midway_app ├── src │ ├── consumer │ │ └── user.consumer.ts │ ├── interface.ts │ └── service │ └── user.service.ts ├── test ├── package.json └── tsconfig.json最小消费者示例import { Consumer, MSListenerType, RabbitMQListener, Inject } from midwayjs/core; import { Context } from midwayjs/rabbitmq; import { ConsumeMessage } from amqplib; Consumer(MSListenerType.RABBITMQ) export class UserConsumer { Inject() ctx: Context; RabbitMQListener(tasks) async gotData(msg: ConsumeMessage) { this.ctx.channel.ack(msg); } }Consumer(MSListenerType.RABBITMQ)声明这是一个 RabbitMQ 类型的消费者。MSListenerType枚举同时覆盖 MQTT、Kafka 等类型见 consumer.ts框架据此按类型装载订阅者。RabbitMQListener(tasks)第一个参数为要监听的队列名将当前方法绑定到名为tasks的队列。订阅多个队列只需写多个方法或拆多个文件。方法入参msg是amqplib的ConsumeMessage。默认情况下消息需要显式 ack通过注入的ctx.channel.ack(msg)确认消费成功也可以配置noAck让服务端自动确认。消息上下文Context订阅消息的上下文与 Web 上下文类似内部包含独立的requestContext与每次接收到的消息数据绑定。从框架导出的类型定义为见 interface.tsexport type IMidwayRabbitMQContext IMidwayContext{ data: ConsumeMessage; // 原始消息 channel: Channel; // 当前消费使用的 channel queueName: string; // 队列名 ack: (data: any) void; // 便捷 ack 方法 };消费代码中可直接注入使用import { Context } from midwayjs/rabbitmq;配置详解url、socketOptions 与 reconnectTime在src/config/config.default.ts中配置连接信息// src/config/config.default.ts import { MidwayConfig } from midwayjs/core; export default { rabbitmq: { url: amqp://localhost } } as MidwayConfig;核心配置项类型定义见 interface.ts属性描述默认值/说明urlRabbitMQ 连接信息可以是amqp://localhost这样的字符串也可以是amqplib的Options.Connect对象socketOptionsamqplib.connect的第二个参数即 socket 级选项可选透传给底层连接reconnectTime断线重连的时间间隔默认 10 秒exchanges需要预声明的交换机列表name/type/options可选useConfirmChannel是否使用 Confirm Channel发布确认模式可选配合createChannel(true)使用其中reconnectTime的默认值在 mq.ts 中有明确实现this.reconnectTime options.reconnectTime ?? 10 * 1000;DefaultConfig类型即string | AmqpOptions.Connectinterface.ts也就是说url支持传入完整的连接选项对象包含 hostname、port、credentials 等。装饰器参数详解RabbitMQListenerOptionsRabbitMQListener的第二个参数是一个配置对象完整类型定义位于 rabbitmqListener.tsexport interface RabbitMQListenerOptions { propertyKey?: string; queueName?: string; exchange?: string; // 队列选项 exclusive?: boolean; durable?: boolean; autoDelete?: boolean; messageTtl?: number; expires?: number; deadLetterExchange?: string; deadLetterRoutingKey?: string; maxLength?: number; maxPriority?: number; pattern?: string; // 预取数量 prefetch?: number; // 路由键 routingKey?: string; // 交换机选项 exchangeOptions?: { type?: direct | topic | headers | fanout | match | string; durable?: boolean; internal?: boolean; autoDelete?: boolean; alternateExchange?: string; arguments?: any; }; // 消费选项 consumeOptions?: { consumerTag?: string; noLocal?: boolean; noAck?: boolean; exclusive?: boolean; priority?: number; arguments?: any; }; }参数含义速查队列层durable持久化默认 true、exclusive独占队列、autoDelete无消费者时自动删除、messageTtl消息过期时间 ms、expires队列过期时间、deadLetterExchange/deadLetterRoutingKey死信配置、maxLength队列最大消息数、maxPriority优先级上限。路由层exchange绑定的交换机名、exchangeOptions.type交换机类型direct/topic/headers/fanout/match、routingKey路由键、patternbinding 使用的匹配模式与 routingKey 二选一。消费层prefetch预取数量默认 1、consumeOptions.noAck是否自动确认、consumeOptions.priority消费者优先级。prefetch默认值在 mq.ts 中体现channel.prefetch(listenerOptions.prefetch ?? 1)。队列的durable默认值同理createConsumer内部使用Object.assign({ durable: true }, listenerOptions)作为assertQueue的选项。Exchange 路由策略Fanout 与 DirectAMQP 的核心概念包括 Producer生产者、Broker队列服务器实体含 Exchange/Binding/Queue与 Consumer消费者。消息由 Producer 发布到 ExchangeExchange 依据路由规则将消息投递到绑定好的 QueueConsumer 通过订阅 Queue 消费。Midway 组件在 mq.ts 中会依次完成assertQueue→assertExchange→bindQueue的装配流程。Fanout Exchange广播Fanout 交换机会忽略 RoutingKey将消息广播给所有绑定的队列。下面示例中abc与bcd两个队列绑定同一个名为logs的 fanout 交换机两条消费者方法都会收到消息import { Consumer, MSListenerType, RabbitMQListener, Inject, App } from midwayjs/core; import { Context, Application } from midwayjs/rabbitmq; import { ConsumeMessage } from amqplib; Consumer(MSListenerType.RABBITMQ) export class UserConsumer { App() app: Application; Inject() ctx: Context; Inject() logger; RabbitMQListener(abc, { exchange: logs, exchangeOptions: { type: fanout, durable: false, }, exclusive: true, consumeOptions: { noAck: true, } }) async gotData(msg: ConsumeMessage) { this.logger.info(test output1 , msg.content.toString(utf8)); } RabbitMQListener(bcd, { exchange: logs, exchangeOptions: { type: fanout, durable: false, }, exclusive: true, consumeOptions: { noAck: true, } }) async gotData2(msg: ConsumeMessage) { this.logger.info(test output2 , msg.content.toString(utf8)); } }Direct Exchange按路由键定向过滤Direct 是 RabbitMQ 的默认交换机类型完全根据 RoutingKey 路由绑定队列时指定 RoutingKey发送消息时携带相同的 RoutingKey消息即被路由到对应队列。示例中不写死 Queue Name仅声明路由键import { Consumer, MSListenerType, RabbitMQListener, Inject, App } from midwayjs/core; import { Context, Application } from midwayjs/rabbitmq; import { ConsumeMessage } from amqplib; Consumer(MSListenerType.RABBITMQ) export class UserConsumer { App() app: Application; Inject() ctx: Context; Inject() logger; RabbitMQListener(, { exchange: direct_logs, exchangeOptions: { type: direct, durable: false, }, routingKey: direct_key, exclusive: true, consumeOptions: { noAck: true, } }) async gotData(msg: ConsumeMessage) { // TODO } }direct 类型的消息会根据 routingKey 做定向过滤只有绑定相同路由键的订阅才能收到消息。源码级原理消费者如何被装载与执行连接建立与断线监听framework.ts 的run()方法清晰地展示了启动链路public async run(): Promisevoid { try { // 建立连接 await this.app.connect( this.configurationOptions.url, this.configurationOptions.socketOptions ); // 装载订阅者 await this.loadSubscriber(); this.logger.info(Rabbitmq server start success); } catch (error) { this.app.close(); throw error; } }在 mq.ts 的connect方法中组件通过amqp-connection-manager建立连接并监听三组关键事件connect连接成功logger.info(Message Queue connected!)并 resolve 启动 PromiseconnectFailed连接失败记录错误并 reject——这正是 CHANGELOG 中3.4.0-beta.11throw error when rabbitmq connect fail 的落地实现保证启动失败不会静默error运行期异常如断线记录错误日志配合reconnectTime由连接管理器自动重连——对应2.11.0引入的 reconnection 能力。createConsumer 的 Channel 装配流程createConsumer 是消费能力的核心。它基于连接管理器创建带setup的 channel在建立连接后按顺序执行assertQueue(queueName, { durable: true, ...listenerOptions })声明队列默认持久化若配置了exchangeassertExchange默认类型topic→bindQueue绑定键取routingKey || patternchannel.prefetch(prefetch ?? 1)设置预取数量channel.consume(queueName, listenerCallback, consumeOptions)注册消息回调并把收到的消息交给上层回调。订阅者装载与调用链loadSubscriber()framework.ts的工作分三步通过DecoratorManager.listModule(MS_CONSUMER_KEY)找出所有标记了Consumer(MSListenerType.RABBITMQ)的模块用listPropertyDataFromClass取出每个类上所有RabbitMQListener绑定的方法及配置为每个监听配置调用app.createConsumer在消息到达时构造ctx含data、channel、queueName、ack然后依次执行guard 校验对应3.6.0新增的守卫能力未通过会抛出MidwayInvokeForbiddenError→ 按需从requestContext解析实例 → 经过中间件管道 → 调用业务方法并传入原始消息。关于 ack 有一个值得注意的设计业务方法有返回值时框架会自动 ackframework.ts而示例中手写的this.ctx.channel.ack(msg)适用于无返回值的场景。两者择一即可切勿重复确认。链路追踪Tracing集成从 mq.ts 可以看到组件对 Midway 链路追踪的原生支持创建 channel 后通过bindTraceContext包装sendToQueue与publish方法自动注入 trace 头消费侧在 framework.ts 中用MidwayTraceService.runWithEntrySpan包裹业务调用并透出midway.protocol rabbitmq、midway.rabbitmq.queue等属性。若你的应用开启了 tracing相关配置见 site/docs/tracing.mdRabbitMQ 的收发链路会自动接入。本地测试无需真实 MQ 的 Mock 生产端midwayjs/mock提供了createRabbitMQProducer方法实现见 rabbitMQ.ts它默认以mock: true运行——此时amqplib.connect会被替换为纯内存实现队列、交换机、消息都保存在进程内无需启动 RabbitMQ 服务即可联调。基础用法发消息到队列import { createRabbitMQProducer, close, creatApp } from midwayjs/mock; describe(/test/index.test.ts, () { it(should test create message and get from app, async () { // 创建队列和 channel const channel await createRabbitMQProducer(tasks, { isConfirmChannel: true, mock: false, // 设为 false 则连接真实 amqp://localhost url: amqp://localhost, }); // 向队列发送数据 channel.sendToQueue(tasks, Buffer.from(something to do)); // 启动应用自动监听队列 const app await creatApp(); await close(app); }); });示例一Fanout 交换机测试const manager await createRabbitMQProducer(tasks-fanout, { isConfirmChannel: false, mock: false, url: amqp://localhost, }); const ex logs; const msg Hello World!; // 声明 fanout 交换机广播给所有绑定的队列 manager.assertExchange(ex, fanout, { durable: false }); // 启动服务可通过 reconnectTime 缩短重连等待 const app await creatApp(base-app-fanout, { url: amqp://localhost, reconnectTime: 2000, }); // 发送到交换机非持久化场景需等订阅服务起来后再发 manager.sendToExchange(ex, , Buffer.from(msg)); await sleep(5000); await manager.close(); await close(app);示例二Direct 交换机测试const manager await createRabbitMQProducer(tasks-direct, { isConfirmChannel: false, mock: false, url: amqp://localhost, }); const ex direct_logs; const msg Hello World!; manager.assertExchange(ex, direct, { durable: false }); const app await creatApp(base-app-direct, { url: amqp://localhost, reconnectTime: 2000, }); // 指定 routingKey 发送只有绑定相同 key 的队列能收到 manager.sendToExchange(ex, direct_key, Buffer.from(msg)); await manager.close(); await close(app);Mock 模式的内存实现rabbitMQ.ts对 fanout/direct/headers 等交换机分别做了路由模拟fanout 向所有绑定队列广播direct 则精确匹配pattern routingKey的绑定——与真实 AMQP 语义一致可放心用于断言消费结果。生产者Producer实践纯 SDK 封装Midway 当前没有为消息发送提供组件化封装官方推荐用amqplibamqp-connection-manager直接封装一个全局单例 Service。首先安装依赖$ npm i amqplib amqp-connection-manager --save $ npm i types/amqplib --save-dev然后在src/service/rabbitmq.ts中封装import { Provide, Scope, ScopeEnum, Init, Autoload, Destroy } from midwayjs/core; import * as amqp from amqp-connection-manager; Autoload() Provide() Scope(ScopeEnum.Singleton) // 单例进程内全局唯一 export class RabbitmqService { private connection: amqp.AmqpConnectionManager; private channelWrapper; Init() async connect() { // 创建连接建议把配置放到 Config 中再注入 this.connection await amqp.connect(amqp://localhost); // 创建 channel并在 setup 中声明队列 this.channelWrapper this.connection.createChannel({ json: true, setup: function(channel) { return Promise.all([ channel.assertQueue(tasks, { durable: true }), ]); }, }); } public async sendToQueue(queueName: string, data: any) { return this.channelWrapper.sendToQueue(queueName, data); } Destroy() async close() { await this.channelWrapper.close(); await this.connection.close(); } }通过Autoload使该 Service 随应用自启动Scope(ScopeEnum.Singleton)保证连接全局唯一。业务侧注入后直接调用Provide() export class UserService { Inject() rabbitmqService: RabbitmqService; async invoke() { // 发送消息 await this.rabbitmqService.sendToQueue(tasks, { hello: world }); } }小结midwayjs/rabbitmq是一个小而完整的消息订阅组件它以amqp-connection-manager为底座获得了自动重连能力CHANGELOG2.11.0在启动语义上严格做到连接失败即抛错3.4.0-beta.11并在3.6.0之后让消费者完全接入 Midway 的 guard 体系消费链路从assertQueue/assertExchange/bindQueue的 channel 装配到requestContext内的守卫、中间件与方法调用再到自动/手动 ack 与 tracing 透传每一层都能在源码中找到对应实现。搭配midwayjs/mock的纯内存createRabbitMQProducer开发者无需依赖真实 MQ 实例即可完成端到端联调。延伸阅读组件源码入口与类型定义packages/rabbitmq/src/index.ts、packages/rabbitmq/src/interface.ts连接管理与消费核心实现packages/rabbitmq/src/mq.ts框架装载与订阅者调度packages/rabbitmq/src/framework.ts监听装饰器与选项类型定义packages/core/src/decorator/microservice/rabbitmqListener.ts本地测试 Mock 生产端packages/mock/src/client/rabbitMQ.ts官方使用文档site/docs/extensions/rabbitmq.md赞分享后端微服务云原生【免费下载链接】midway A Node.js Serverless Framework for front-end/full-stack developers. Build the application for next decade. Works on AWS, Alibaba Cloud, Tencent Cloud and traditional VM/Container. Super easy integrate with React and Vue. 项目地址https://gitcode.com/gh_mirrors/mi/midway点击查看免费下载相关推荐Apache Pulsar 消息机制深度解析从生产、消费到订阅模型Apache Pulsar 消息机制深度解析从生产、消费到订阅模型 Apache Pulsar 构建于经典的发布 订阅pub sub模式之上生产者pr消息队列后端流处理svgr/rollup 完整指南从版本演进到源码级的 Rollup SVG 转 React 组件实践svgr/rollup 完整指南从版本演进到源码级的 Rollup SVG 转 React 组件实践 SVGRSVG to React为 Rollup前端开发工具Victory Legend 组件的演进全记录从版本变更史到源码级实现解析Victory Legend 组件的演进全记录从版本变更史到源码级实现解析 Victory 是一套用于构建交互式数据可视化的 React 组件库而 vict数据可视化UI组件上一篇如何通过iPXE实现企业级网络引导的现代化升级下一篇终极指南Koel音乐流媒体平台的用户权限管理与API安全机制解析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考