消息队列是现代分布式系统的神经中枢。在 Node.js 生态中,从任务调度到微服务通信,从事件溯源到流式处理,消息队列都扮演着不可替代的角色。
1. 消息队列核心概念
1.1 为什么需要消息队列
在高并发系统中,同步调用会导致级联故障。消息队列通过异步解耦将生产者和消费者分离,使系统具备弹性的伸缩能力。
同步调用(紧耦合):
API → 调用订单服务 → 调用库存服务 → 调用支付服务
任一环节失败 → 整个链路回滚 → 用户体验差
异步消息(松耦合):
API → 发送订单消息 ──┬─→ 订单服务
├─→ 库存服务
├─→ 支付服务
└─→ 通知服务
各服务独立消费 → 单点故障不影响全局
1.2 三大核心能力
| 能力 | 说明 | 典型场景 |
|---|---|---|
| 异步解耦 | 生产者与消费者独立演进 | 订单创建后异步通知物流 |
| 削峰填谷 | 缓冲突发流量,平滑处理 | 秒杀活动峰值流量 |
| 可靠投递 | 消息持久化 + 重试机制 | 金融交易通知 |
1.3 背压(Backpressure)
当消费者处理速度远低于生产者时,队列深度会无限增长,最终耗尽内存。背压机制通过限流、阻塞、丢弃三种策略防止系统过载。
// Node.js Stream 中的背压处理
const readable = getMessageStream();
const writable = getConsumerStream();
readable.on('data', (chunk) => {
const canContinue = writable.write(chunk);
if (!canContinue) {
readable.pause(); // 暂停生产
writable.once('drain', () => {
readable.resume(); // 恢复生产
});
}
});
1.4 重试与死信
消息处理失败后不应立即丢弃,而应进入重试队列。当重试次数耗尽,消息进入**死信队列(DLQ)**供人工排查。
消息生命周期:
正常队列 → 消费失败 → 重试队列(延迟 5s)→ 再次消费
↓
再次失败 → 重试队列(延迟 30s)
↓
再次失败 → 死信队列(人工处理)
2. RabbitMQ 与 amqplib
RabbitMQ 是最成熟的开源消息代理,基于 AMQP 协议,支持复杂的路由规则和可靠投递。
2.1 核心概念
┌─────────────┐ ┌──────────┐ ┌───────────┐ ┌────────────┐
│ Producer │────→│ Exchange │────→│ Queue │────→│ Consumer │
└─────────────┘ └──────────┘ └───────────┘ └────────────┘
│ ↑
└──── Routing Key 匹配规则 ──────┘
| 概念 | 说明 |
|---|---|
| Exchange | 消息路由器,决定消息进入哪个队列 |
| Queue | 消息存储的缓冲区 |
| Binding | Exchange 与 Queue 之间的绑定关系 |
| Routing Key | 消息的路由标识,Exchange 依据它进行分发 |
2.2 Exchange 类型
const amqp = require('amqplib');
// 连接到 RabbitMQ
const conn = await amqp.connect('amqp://guest:guest@localhost:5672');
const ch = await conn.createChannel();
// 1. Direct Exchange:精确匹配 Routing Key
await ch.assertExchange('orders.direct', 'direct', { durable: true });
await ch.assertQueue('orders.payment');
await ch.bindQueue('orders.payment', 'orders.direct', 'payment');
// 2. Topic Exchange:模式匹配(支持 * 和 #)
await ch.assertExchange('logs.topic', 'topic', { durable: true });
await ch.assertQueue('logs.error');
await ch.bindQueue('logs.error', 'logs.topic', 'kernel.error.*');
// 3. Fanout Exchange:广播到所有绑定队列
await ch.assertExchange('notifications.fanout', 'fanout', { durable: true });
// 4. Headers Exchange:根据消息头属性匹配
await ch.assertExchange('events.headers', 'headers', { durable: true });
2.3 生产者与消费者完整示例
// producer.js
const amqp = require('amqplib');
async function sendOrder(order) {
const conn = await amqp.connect(process.env.RABBITMQ_URL);
const ch = await conn.createChannel();
// 声明 durable Exchange 和 Queue
await ch.assertExchange('orders.topic', 'topic', { durable: true });
await ch.assertQueue('orders.processing', {
durable: true,
arguments: {
'x-dead-letter-exchange': 'orders.dlx',
'x-dead-letter-routing-key': 'orders.failed'
}
});
await ch.bindQueue('orders.processing', 'orders.topic', 'order.created');
// 发送消息(persistent 确保持久化到磁盘)
const msg = Buffer.from(JSON.stringify(order));
ch.publish('orders.topic', 'order.created', msg, {
persistent: true,
messageId: order.id,
timestamp: Date.now()
});
console.log(`[Producer] Order ${order.id} sent`);
await ch.close();
await conn.close();
}
// consumer.js
async function consumeOrders() {
const conn = await amqp.connect(process.env.RABBITMQ_URL);
const ch = await conn.createChannel();
// 每次只接收 1 条,处理完再取下一条
await ch.prefetch(1);
await ch.consume('orders.processing', async (msg) => {
if (!msg) return;
try {
const order = JSON.parse(msg.content.toString());
console.log(`[Consumer] Processing order ${order.id}`);
// 模拟业务处理
await processOrder(order);
// 确认消息已处理
ch.ack(msg);
} catch (err) {
console.error(`[Consumer] Failed:`, err.message);
// 拒绝消息,requeue=false 则进入死信队列
ch.nack(msg, false, false);
}
});
}
async function processOrder(order) {
// 订单处理逻辑
await new Promise(resolve => setTimeout(resolve, 100));
}
2.4 死信队列(DLQ)配置
// 声明死信交换器和队列
await ch.assertExchange('orders.dlx', 'topic', { durable: true });
await ch.assertQueue('orders.dead', { durable: true });
await ch.bindQueue('orders.dead', 'orders.dlx', 'orders.failed');
// 主队列绑定死信参数(在 assertQueue 时声明)
await ch.assertQueue('orders.processing', {
durable: true,
arguments: {
'x-message-ttl': 300000, // 消息 5 分钟过期
'x-dead-letter-exchange': 'orders.dlx',
'x-dead-letter-routing-key': 'orders.failed',
'x-max-retries': 3 // 最大重试次数(配合插件或代码实现)
}
});
3. Apache Kafka 与 KafkaJS
Kafka 是分布式流处理平台,以高吞吐量和持久化日志著称,适合大数据量、高并发的实时流场景。
3.1 Kafka 核心概念
┌──────────────────────────────────────────────┐
│ Kafka Cluster │
│ ┌─────────────┐ ┌─────────────┐ │
│ │ Broker 1 │ │ Broker 2 │ │
│ └─────────────┘ └─────────────┘ │
│ │ │ │
│ ┌─────────────────────────────────────────┐│
│ │ Topic: "orders" ││
│ │ ┌─────────┐ ┌─────────┐ ┌─────────┐ ││
│ │ │ P0 (Leader)││ P1 │ │ P2 │ ││
│ │ │ R0 │ │ R0 │ │ R0 │ ││
│ │ │ R1 │ │ R1 │ │ R1 │ ││
│ │ └─────────┘ └─────────┘ └─────────┘ ││
│ └─────────────────────────────────────────┘│
│ ↑ ↑ │
│ ┌─────────────┐ ┌─────────────┐ │
│ │ Consumer G1 │ │ Consumer G1 │ │
│ │ (Instance1)│ │ (Instance2)│ │
│ └─────────────┘ └─────────────┘ │
└──────────────────────────────────────────────┘
P = Partition, R = Replica
| 概念 | 说明 |
|---|---|
| Topic | 消息主题,逻辑上的消息分类 |
| Partition | Topic 的分片,每个 Partition 是有序的日志序列 |
| Offset | 消息在 Partition 中的位置标识 |
| Consumer Group | 消费者组,组内消费者共同消费一个 Topic |
| Replication | 分区副本,保证高可用 |
3.2 KafkaJS 生产者与消费者
const { Kafka } = require('kafkajs');
const kafka = new Kafka({
clientId: 'order-service',
brokers: ['kafka1:9092', 'kafka2:9092'],
retry: {
initialRetryTime: 100,
retries: 8
}
});
// ========== 生产者 ==========
const producer = kafka.producer({
idempotent: true, // 幂等生产者,防止重复发送
transactionalId: 'order-producer'
});
async function sendOrderEvent(order) {
await producer.connect();
// 发送带 Key 的消息(相同 Key 进入同一 Partition,保证顺序)
await producer.send({
topic: 'orders',
messages: [{
key: order.userId, // 按用户 ID 分区
value: JSON.stringify(order),
headers: {
'event-type': 'order.created',
'version': '1.0'
}
}]
});
await producer.disconnect();
}
// ========== 消费者 ==========
const consumer = kafka.consumer({
groupId: 'order-processor-group',
sessionTimeout: 30000,
heartbeatInterval: 3000
});
async function consumeOrderEvents() {
await consumer.connect();
await consumer.subscribe({ topic: 'orders', fromBeginning: false });
await consumer.run({
autoCommit: false, // 手动提交偏移量
eachBatchAutoResolve: false,
eachBatch: async ({ batch, resolveOffset, heartbeat, commitOffsetsIfNecessary }) => {
for (const message of batch.messages) {
try {
const order = JSON.parse(message.value.toString());
console.log(`[Kafka] Processing order: ${order.id}, partition: ${batch.partition}, offset: ${message.offset}`);
await processOrder(order);
// 处理成功,提交偏移量
await resolveOffset(message.offset);
await heartbeat();
} catch (err) {
console.error(`[Kafka] Processing failed:`, err);
// 不提交偏移量,下次重试
throw err;
}
}
await commitOffsetsIfNecessary();
}
});
}
3.3 消费者组重平衡
消费者组内的实例数应与分区数匹配或成比例,过多消费者会导致空闲,过少则消费延迟。
// 监听重平衡事件
consumer.on(consumer.events.GROUP_JOIN, (event) => {
console.log(`Joined group: ${event.payload.groupId}, generation: ${event.payload.generationId}`);
});
consumer.on(consumer.events.REBALANCING, (event) => {
console.log('Consumer group rebalancing...');
});
4. BullMQ 与 Redis 任务队列
BullMQ 是基于 Redis 的 Node.js 队列库,专注任务调度场景,支持延迟任务、重复任务和优先级队列。
4.1 基础队列使用
const { Queue, Worker } = require('bullmq');
const Redis = require('ioredis');
const connection = new Redis({ host: 'localhost', port: 6379, maxRetriesPerRequest: null });
// 创建队列
const emailQueue = new Queue('email', { connection });
// 添加任务
async function enqueueEmail(data) {
const job = await emailQueue.add('send-email', data, {
attempts: 3, // 失败重试 3 次
backoff: {
type: 'exponential', // 指数退避
delay: 5000 // 初始延迟 5 秒
},
removeOnComplete: 100, // 保留最近 100 条完成记录
removeOnFail: 50 // 保留最近 50 条失败记录
});
console.log(`Job ${job.id} enqueued`);
}
// 创建 Worker 处理任务
const emailWorker = new Worker('email', async (job) => {
console.log(`[Worker] Processing job ${job.id}:`, job.data);
// 模拟发送邮件
await sendEmail(job.data);
return { sent: true, timestamp: Date.now() };
}, {
connection,
concurrency: 5, // 并发处理 5 个任务
limiter: {
max: 100, // 每秒最多 100 个任务
duration: 1000
}
});
emailWorker.on('completed', (job, result) => {
console.log(`Job ${job.id} completed:`, result);
});
emailWorker.on('failed', (job, err) => {
console.error(`Job ${job.id} failed:`, err.message);
});
4.2 延迟任务与重复任务
// ========== 延迟任务 ==========
// 10 秒后执行
await emailQueue.add('scheduled-email', { to: 'user@example.com' }, {
delay: 10000
});
// 指定确切时间执行
const delay = new Date('2026-08-18T09:00:00+08:00').getTime() - Date.now();
await emailQueue.add('morning-digest', { userId: 123 }, { delay });
// ========== 重复任务(Cron 风格)==========
const { QueueScheduler } = require('bullmq');
new QueueScheduler('reports', { connection });
await emailQueue.add('daily-report', { type: 'analytics' }, {
repeat: {
cron: '0 9 * * *', // 每天上午 9 点
tz: 'Asia/Shanghai'
}
});
await emailQueue.add('weekly-summary', {}, {
repeat: {
every: 7 * 24 * 60 * 60 * 1000, // 每 7 天
limit: 52 // 最多执行 52 次
}
});
4.3 队列事件监控
// 监听队列事件
emailQueue.on('waiting', (jobId) => console.log(`Job ${jobId} is waiting`));
emailQueue.on('active', (job) => console.log(`Job ${job.id} is active`));
emailQueue.on('stalled', (jobId) => console.warn(`Job ${jobId} stalled`));
// 获取队列状态
async function getQueueStatus() {
const [waiting, active, completed, failed, delayed] = await Promise.all([
emailQueue.getWaitingCount(),
emailQueue.getActiveCount(),
emailQueue.getCompletedCount(),
emailQueue.getFailedCount(),
emailQueue.getDelayedCount()
]);
return { waiting, active, completed, failed, delayed };
}
5. AWS SQS 与 SNS 集成
AWS 提供托管消息服务,适合云服务原生场景,无需自行运维消息基础设施。
5.1 SQS 标准队列与 FIFO 队列
const { SQSClient, SendMessageCommand, ReceiveMessageCommand, DeleteMessageCommand } = require('@aws-sdk/client-sqs');
const sqs = new SQSClient({ region: 'ap-southeast-1' });
const QUEUE_URL = process.env.SQS_QUEUE_URL;
// 发送消息
async function sendSQSMessage(payload) {
await sqs.send(new SendMessageCommand({
QueueUrl: QUEUE_URL,
MessageBody: JSON.stringify(payload),
MessageAttributes: {
'EventType': { StringValue: 'OrderCreated', DataType: 'String' }
}
}));
}
// 接收并处理消息
async function pollSQS() {
const result = await sqs.send(new ReceiveMessageCommand({
QueueUrl: QUEUE_URL,
MaxNumberOfMessages: 10,
WaitTimeSeconds: 20, // 长轮询
VisibilityTimeout: 300 // 处理超时 5 分钟
}));
for (const message of result.Messages || []) {
try {
const body = JSON.parse(message.Body);
await processMessage(body);
// 删除已处理消息
await sqs.send(new DeleteMessageCommand({
QueueUrl: QUEUE_URL,
ReceiptHandle: message.ReceiptHandle
}));
} catch (err) {
console.error('SQS processing failed:', err);
// 不删除消息,VisibilityTimeout 到期后重试
}
}
}
5.2 SNS + SQS 发布订阅模式
const { SNSClient, PublishCommand } = require('@aws-sdk/client-sns');
const sns = new SNSClient({ region: 'ap-southeast-1' });
const TOPIC_ARN = process.env.SNS_TOPIC_ARN;
// SNS 发布消息
async function publishEvent(event) {
await sns.send(new PublishCommand({
TopicArn: TOPIC_ARN,
Message: JSON.stringify(event),
MessageAttributes: {
'env': { DataType: 'String', StringValue: 'production' }
}
}));
}
| 特性 | SQS 标准队列 | SQS FIFO | SNS |
|---|---|---|---|
| 顺序保证 | 否 | 是(同一 Group) | 否 |
| 去重 | 否 | 是(5 分钟窗口) | 否 |
| 推送/拉取 | 拉取 | 拉取 | 推送 |
| 模式 | 点对点 | 点对点 | 发布订阅 |
| 吞吐量 | 近乎无限 | 3000 TPS | 近乎无限 |
6. 事件驱动架构模式
6.1 发布订阅模式
┌──────────────┐
│ Publisher │
└──────┬───────┘
│ publish(event)
▼
┌──────────────┐
│ Event Bus │
└──────┬───────┘
│
┌─────────┼─────────┐
▼ ▼ ▼
┌─────────┐ ┌─────────┐ ┌─────────┐
│ HandlerA│ │ HandlerB│ │ HandlerC│
│ Email │ │ Analytics│ │ Cache │
└─────────┘ └─────────┘ └─────────┘
6.2 事件溯源(Event Sourcing)
状态变更不以当前状态存储,而是以事件流的形式追加记录。通过重放事件可重建任意时刻的状态。
// 事件存储
const events = [
{ type: 'OrderCreated', data: { orderId: '001', items: [...] } },
{ type: 'PaymentProcessed', data: { orderId: '001', amount: 199 } },
{ type: 'OrderShipped', data: { orderId: '001', tracking: 'SF123' } }
];
// 状态重建
function rebuildState(events) {
return events.reduce((state, event) => {
switch (event.type) {
case 'OrderCreated':
return { ...state, id: event.data.orderId, items: event.data.items, status: 'created' };
case 'PaymentProcessed':
return { ...state, status: 'paid', paidAt: Date.now() };
case 'OrderShipped':
return { ...state, status: 'shipped', tracking: event.data.tracking };
default:
return state;
}
}, {});
}
6.3 CQRS 与事件总线
命令查询职责分离(CQRS)配合事件总线,将写操作和读操作分离到不同模型,通过事件同步状态。
写模型(Command Side) 事件总线 读模型(Query Side)
┌──────────────┐ ┌──────────┐ ┌──────────────┐
│ 创建订单命令 │─→ 订单聚合根 ─│→ OrderCreated│→─ │ 订单视图更新 │
└──────────────┘ │→ PaymentDone │→─ │ 搜索索引更新 │
└──────────┘ └──────────────┘
7. Saga 模式实现
分布式事务无法依赖传统的 ACID,Saga 模式通过补偿事务保证最终一致性。
7.1 编排式 Saga(Choreography)
每个服务完成本地事务后广播事件,触发下一个服务的动作。
OrderService ──→ OrderCreated ──→ InventoryService ──→ StockReserved
│
▼
PaymentService ←── PaymentRequest ←── OrderService
│
▼
PaymentProcessed ──→ ShipmentService ──→ OrderCompleted
7.2 编排式 Saga 代码示例
// order-service/saga.js
const { Kafka } = require('kafkajs');
const kafka = new Kafka({ clientId: 'order-saga', brokers: ['kafka:9092'] });
const producer = kafka.producer();
const consumer = kafka.consumer({ groupId: 'order-saga-group' });
// Saga 步骤定义
const sagaSteps = {
'order.created': {
next: 'inventory.reserve',
compensate: 'order.cancel'
},
'inventory.reserved': {
next: 'payment.charge',
compensate: 'inventory.release'
},
'payment.success': {
next: 'shipment.create',
compensate: 'payment.refund'
},
'shipment.created': {
next: null, // Saga 完成
compensate: 'shipment.cancel'
}
};
async function startSaga(order) {
await producer.connect();
await producer.send({
topic: 'saga-events',
messages: [{
key: order.id,
value: JSON.stringify({
type: 'order.created',
payload: order,
sagaId: order.id,
step: 0
})
}]
});
}
async function handleSagaEvent(event) {
const { type, payload, sagaId } = event;
const step = sagaSteps[type];
if (!step) return;
try {
// 执行本地事务
await executeLocalTransaction(type, payload);
if (step.next) {
// 触发下一步
await producer.send({
topic: 'saga-events',
messages: [{
key: sagaId,
value: JSON.stringify({
type: step.next,
payload,
sagaId,
step: event.step + 1
})
}]
});
}
} catch (err) {
// 触发补偿事务
console.error(`Saga step ${type} failed, triggering compensation`);
await triggerCompensation(sagaId, step.compensate, payload);
}
}
async function triggerCompensation(sagaId, compensateAction, payload) {
await producer.send({
topic: 'saga-compensation',
messages: [{
key: sagaId,
value: JSON.stringify({ type: compensateAction, payload, sagaId })
}]
});
}
7.3 协调式 Saga(Orchestration)
由专门的 Saga 编排器集中管理事务流程,适合复杂业务流程。
// saga-orchestrator.js
class OrderSagaOrchestrator {
constructor() {
this.state = new Map(); // sagaId -> currentStep
}
async execute(orderId, orderData) {
const saga = { id: orderId, status: 'pending', steps: [] };
try {
// Step 1: 创建订单
await this.callService('order-service', 'create', orderData);
saga.steps.push({ service: 'order', action: 'create', status: 'ok' });
// Step 2: 预留库存
await this.callService('inventory-service', 'reserve', { orderId, items: orderData.items });
saga.steps.push({ service: 'inventory', action: 'reserve', status: 'ok' });
// Step 3: 扣款
await this.callService('payment-service', 'charge', { orderId, amount: orderData.total });
saga.steps.push({ service: 'payment', action: 'charge', status: 'ok' });
// Step 4: 创建物流
await this.callService('shipment-service', 'create', { orderId, address: orderData.address });
saga.status = 'completed';
} catch (err) {
saga.status = 'failed';
// 倒序执行补偿
for (const step of [...saga.steps].reverse()) {
await this.compensate(step, orderId);
}
}
return saga;
}
async compensate(step, orderId) {
const compensations = {
'order:create': () => this.callService('order-service', 'cancel', { orderId }),
'inventory:reserve': () => this.callService('inventory-service', 'release', { orderId }),
'payment:charge': () => this.callService('payment-service', 'refund', { orderId })
};
const key = `${step.service}:${step.action}`;
if (compensations[key]) {
await compensations[key]();
}
}
}
8. 消息顺序与 Exactly-Once 语义
8.1 顺序保证策略
| 场景 | 方案 | 说明 |
|---|---|---|
| 全局顺序 | 单分区/单队列 | 吞吐量受限 |
| 分区顺序 | 按 Key 分区 | 同一 Key 的消息进入同一分区 |
| 因果顺序 | 向量时钟 | 分布式系统中事件因果关系 |
// Kafka 按 Key 分区保证用户级顺序
await producer.send({
topic: 'user-events',
messages: userEvents.map(e => ({
key: e.userId, // 同一用户的事件进入同一分区
value: JSON.stringify(e)
}))
});
// RabbitMQ 单队列保证顺序(取消并发消费)
await ch.prefetch(1); // 一次只处理一条
8.2 Exactly-Once 语义
消息消费语义有三种级别:
At-Most-Once: 消息可能丢失,但不会重复
At-Least-Once: 消息不会丢失,但可能重复
Exactly-Once: 消息既不丢失也不重复(理想状态)
实现 Exactly-Once 的关键技术:
// 幂等消费者设计
const processedIds = new Set(); // 生产环境使用 Redis/DB
async function processMessage(msg) {
const id = msg.messageId || msg.id;
// 检查是否已处理
if (await isProcessed(id)) {
console.log(`Message ${id} already processed, skipping`);
return;
}
// 业务处理 + 记录处理状态应在同一事务中
await withTransaction(async (trx) => {
await executeBusinessLogic(msg, trx);
await markAsProcessed(id, trx);
});
}
// Kafka 事务型 Exactly-Once
const transaction = await producer.transaction();
try {
await transaction.send({ topic: 'output', messages: [...] });
await transaction.sendOffsets({
consumerGroupId: 'my-group',
topics: [{
topic: 'input',
partitions: [{ partition: 0, offset: '10' }]
}]
});
await transaction.commit();
} catch (err) {
await transaction.abort();
}
9. 可观测性:监控与告警
9.1 关键监控指标
| 指标 | 说明 | 告警阈值 |
|---|---|---|
| 队列深度 | 未消费消息数量 | > 10000 |
| 消费延迟(Lag) | 消息从生产到消费的间隔 | > 30s |
| 死信数量 | 进入死信队列的消息数 | > 0 |
| 重试率 | 消费失败重试的比例 | > 5% |
| 吞吐率 | 每秒处理消息数 | 基线 +-20% |
9.2 RabbitMQ 监控
// 使用 rabbitmq-management HTTP API
const axios = require('axios');
async function getRabbitMQMetrics() {
const res = await axios.get('http://localhost:15672/api/queues', {
auth: { username: 'guest', password: 'guest' }
});
return res.data.map(q => ({
name: q.name,
messages_ready: q.messages_ready, // 待消费消息
messages_unacknowledged: q.messages_unacknowledged, // 已投递未确认
consumers: q.consumers,
message_stats: q.message_stats
}));
}
9.3 Kafka 消费延迟监控
const { Kafka } = require('kafkajs');
const kafka = new Kafka({ clientId: 'monitor', brokers: ['kafka:9092'] });
const admin = kafka.admin();
async function getConsumerLag(groupId) {
await admin.connect();
const lag = await admin.fetchOffsets({ groupId, topic: 'orders' });
const offsets = await admin.fetchTopicOffsets('orders');
const result = lag.map(({ partition, offset }) => {
const latest = offsets.find(o => o.partition === partition);
return {
partition,
consumerOffset: parseInt(offset),
logEndOffset: latest.offset,
lag: latest.offset - parseInt(offset)
};
});
await admin.disconnect();
return result;
}
9.4 BullMQ 监控面板
const { Queue } = require('bullmq');
const queue = new Queue('email', { connection });
async function getBullMetrics() {
const [waiting, active, completed, failed, delayed, paused] = await Promise.all([
queue.getWaitingCount(),
queue.getActiveCount(),
queue.getCompletedCount(),
queue.getFailedCount(),
queue.getDelayedCount(),
queue.getPausedCount()
]);
const metrics = {
waiting, active, completed, failed, delayed, paused,
total: waiting + active + completed + failed + delayed
};
// 推送到 Prometheus
metricsQueueDepth.set({ queue: 'email' }, waiting);
metricsActiveJobs.set({ queue: 'email' }, active);
metricsFailedJobs.inc({ queue: 'email' }, failed);
return metrics;
}
9.5 健康检查集成
// Express 健康检查端点
app.get('/health/queue', async (req, res) => {
const statuses = await Promise.all([
checkRabbitMQConnection(),
checkKafkaConnection(),
checkRedisConnection()
]);
const allHealthy = statuses.every(s => s.healthy);
res.status(allHealthy ? 200 : 503).json({
status: allHealthy ? 'healthy' : 'unhealthy',
services: statuses
});
});
10. 选型对比与决策指南
10.1 综合对比
| 维度 | RabbitMQ | Kafka | BullMQ | AWS SQS |
|---|---|---|---|---|
| 协议 | AMQP | 自定义二进制协议 | Redis 协议 | HTTP/JSON |
| 模型 | 消息代理 | 分布式日志 | 任务队列 | 托管队列 |
| 顺序保证 | 单队列内有序 | 分区内有序 | FIFO 模式 | FIFO 队列 |
| 消息持久化 | 磁盘/内存 | 磁盘(高可靠) | Redis 持久化 | AWS 托管 |
| 吞吐量 | 万级/秒 | 百万级/秒 | 万级/秒 | 万级/秒 |
| 延迟 | 毫秒级 | 毫秒级 | 毫秒级 | 秒级(可能) |
| 重试机制 | 插件/代码 | 代码实现 | 内置 | 内置 |
| 死信队列 | 原生支持 | 需代码实现 | 原生支持 | DLQ 原生 |
| 延迟任务 | 插件 | 不好支持 | 原生 | 不支持 |
| 运维成本 | 中等 | 高 | 低 | 无 |
| 适用场景 | 交易、路由 | 大数据流 | 任务调度 | 云原生事件 |
10.2 选型决策树
是否需要任务调度(延迟/重复/Cron)?
├─ 是 → BullMQ(基于 Redis)
└─ 否 → 继续问...
是否在 AWS 生态且无自运维能力?
├─ 是 → AWS SQS/SNS
└─ 否 → 继续问...
是否需要极高吞吐(>10万/秒)和流式处理?
├─ 是 → Kafka
└─ 否 → RabbitMQ
10.3 混合架构策略
实际生产中往往是混合使用:
┌────────────────────────────────────────────┐
│ 网关/入口层 │
└──────────────┬─────────────────────────────┘
│
┌──────┴──────┐
▼ ▼
┌──────────────┐ ┌──────────────┐
│ Kafka │ │ BullMQ │
│ 实时事件流 │ │ 任务队列 │
│ 数据管道 │ │ 发送邮件 │
└──────┬───────┘ │ 定时任务 │
│ └──────────────┘
▼
┌──────────────┐
│ RabbitMQ │
│ 微服务通信 │
│ 复杂路由 │
└──────────────┘
| 服务 | 用途 |
|---|---|
| Kafka | 用户行为日志、实时数据分析、事件溯源 |
| RabbitMQ | 订单状态流转、支付回调、跨服务通知 |
| BullMQ | 邮件发送、报表生成、定时数据清理 |
| SQS | 与 Lambda 集成的无服务器事件处理 |
延伸阅读
总结
消息队列是 Node.js 从单进程应用走向分布式系统的关键桥梁。选择合适的消息中间件需要综合考量:
- RabbitMQ 适合需要复杂路由规则和可靠事务投递的业务场景
- Kafka 适合大数据量、高吞吐的实时流处理和事件溯源
- BullMQ 适合任务调度、延迟执行和优先级队列场景
- AWS SQS/SNS 适合云原生、无服务器架构中的事件驱动集成
掌握消息队列的核心原理——异步解耦、背压控制、重试策略、顺序语义和可观测性——是构建高可用分布式系统的必修课。结合 Saga 模式处理分布式事务,配合事件驱动架构实现系统解耦,能让 Node.js 服务在复杂业务环境中保持弹性和可扩展性。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。