Node.js 消息队列集成实战:RabbitMQ、Kafka、BullMQ 与事件驱动架构

全面解析 Node.js 消息队列技术:RabbitMQ、Kafka、BullMQ/Bull、AWS SQS/SNS 的集成实践,涵盖异步解耦、背压控制、死信队列、Saga 模式、Exactly-Once 语义与可观测性监控。

消息队列是现代分布式系统的神经中枢。在 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消息存储的缓冲区
BindingExchange 与 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消息主题,逻辑上的消息分类
PartitionTopic 的分片,每个 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 FIFOSNS
顺序保证是(同一 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 综合对比

维度RabbitMQKafkaBullMQAWS 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 从单进程应用走向分布式系统的关键桥梁。选择合适的消息中间件需要综合考量:

  1. RabbitMQ 适合需要复杂路由规则和可靠事务投递的业务场景
  2. Kafka 适合大数据量、高吞吐的实时流处理和事件溯源
  3. BullMQ 适合任务调度、延迟执行和优先级队列场景
  4. AWS SQS/SNS 适合云原生、无服务器架构中的事件驱动集成

掌握消息队列的核心原理——异步解耦、背压控制、重试策略、顺序语义和可观测性——是构建高可用分布式系统的必修课。结合 Saga 模式处理分布式事务,配合事件驱动架构实现系统解耦,能让 Node.js 服务在复杂业务环境中保持弹性和可扩展性。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「nodejs」更多文章

  1. Node.js ORM 深度对比:Prisma、TypeORM、Sequelize 与 Drizzle
  2. Node.js 设计模式与最佳实践:从 SOLID 到六边形架构
  3. Node.js 高级测试策略:从单元测试到混沌工程的完整实践