You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

NodeJS+RabbitMQ服务重启时如何确保消息被正确处理?

Graceful Shutdown for RabbitMQ Consumer in Node.js

Great question—handling graceful shutdown for message consumers is crucial to avoid lost work or partial processing. Let's break down how to implement the "finish current message then stop" behavior you're looking for, with adjustments to your existing code.

First, let's address the gaps in your current setup:

  • Using noAck: true means RabbitMQ marks messages as delivered immediately. If your script restarts right after receiving a message (before spawning the child process), that message is lost forever.
  • There's no logic to handle shutdown signals, so the process can exit mid-message-processing, leaving your php command incomplete.

Key Changes Needed

  1. Enable manual message acknowledgments: Replace noAck: true with noAck: false, and only acknowledge messages after your child process completes successfully. This ensures RabbitMQ requeues the message if your process crashes before finishing.
  2. Track shutdown state: Add a flag to signal when the process should stop accepting new messages.
  3. Listen for shutdown signals: Catch SIGINT (Ctrl+C) and SIGTERM (system stop commands) to trigger graceful shutdown.
  4. Clean up resources: Once no more messages are being processed, close the RabbitMQ channel and connection before exiting the process.

Modified Code

const amqp = require('amqplib/callback_api');
const { spawn } = require('child_process');
const path = require('path'); // Don't forget to import path!

let isShuttingDown = false;
let consumerTag = null;
let channel = null;
let conn = null;

amqp.connect('amqp://guest:guest@127.0.0.1:5672', (err, connection) => {
  if (err) {
    console.error('Failed to connect to RabbitMQ:', err);
    process.exit(1);
  }
  conn = connection;

  conn.createChannel((err, ch) => {
    if (err) {
      console.error('Failed to create channel:', err);
      conn.close();
      process.exit(1);
    }
    channel = ch;
    const q = 'symfony_messages';

    channel.assertQueue(q, { durable: false });
    console.log(" [*] Waiting for messages in %s. To exit press CTRL+C", q);

    // Start consuming messages, save the consumer tag
    channel.consume(q, (msg) => {
      // If we're shutting down, reject the message so RabbitMQ requeues it
      if (isShuttingDown) {
        channel.nack(msg, false, true);
        return;
      }

      try {
        const event = JSON.parse(msg.content.toString());
        if (event.name === 'product.created') {
          console.log('Indexing order for product ID:', event.payload.product_id);
          const cp = spawn('php', [
            path.join(__dirname, '..', '..', 'bin', 'console'),
            'elastic:index:orders',
            event.payload.product_id
          ]);

          cp.stdout.on('data', (data) => {
            console.log(`stdout: ${data}`);
          });

          cp.stderr.on('data', (data) => {
            console.error(`stderr: ${data}`);
          });

          cp.on('close', (code) => {
            console.log(`Child process exited with code ${code}`);
            // Acknowledge the message only after the child process finishes
            channel.ack(msg);

            // If we're shutting down and no more messages are processing, clean up
            if (isShuttingDown) {
              closeResources();
            }
          });

          cp.on('error', (err) => {
            console.error('Failed to spawn child process:', err);
            // Reject the message so it gets requeued
            channel.nack(msg, false, true);

            if (isShuttingDown) {
              closeResources();
            }
          });
        } else {
          // Acknowledge non-relevant messages immediately
          channel.ack(msg);
        }
      } catch (err) {
        console.error('Failed to process message:', err);
        channel.nack(msg, false, true); // Requeue on parse error

        if (isShuttingDown) {
          closeResources();
        }
      }
    }, { noAck: false }, (err, ok) => {
      if (!err) {
        consumerTag = ok.consumerTag;
      }
    });
  });

  // Handle shutdown signals
  process.on('SIGINT', initiateShutdown);
  process.on('SIGTERM', initiateShutdown);
});

function initiateShutdown() {
  if (isShuttingDown) return; // Avoid duplicate shutdowns
  console.log('\nInitiating graceful shutdown...');
  isShuttingDown = true;

  // Stop receiving new messages from RabbitMQ
  if (channel && consumerTag) {
    channel.cancel(consumerTag)
      .then(() => console.log('Stopped consuming new messages'))
      .catch(err => console.error('Failed to cancel consumer:', err));
  }

  // If no messages are being processed, close immediately
  closeResources();
}

function closeResources() {
  // Wait a short time to ensure in-progress messages finish (optional but safe)
  setTimeout(() => {
    if (channel) {
      channel.close()
        .then(() => console.log('Channel closed'))
        .catch(err => console.error('Failed to close channel:', err));
    }
    if (conn) {
      conn.close()
        .then(() => {
          console.log('Connection closed');
          process.exit(0);
        })
        .catch(err => {
          console.error('Failed to close connection:', err);
          process.exit(1);
        });
    }
  }, 1000);
}

How It Works

  • Manual Acknowledgments: We only call channel.ack(msg) after the child process exits successfully. If the process crashes mid-execution, channel.nack(msg, false, true) tells RabbitMQ to requeue the message so it can be processed later.
  • Shutdown Flag: isShuttingDown prevents new messages from being processed once a shutdown signal is received. Any messages received after shutdown starts are rejected and requeued.
  • Signal Handling: When SIGINT or SIGTERM is received, we cancel the consumer to stop new messages, then wait for in-progress child processes to finish before closing the channel and connection.
  • Resource Cleanup: The closeResources function ensures we properly shut down RabbitMQ connections to avoid leaks.

This setup ensures your script will finish processing the current message (including waiting for the php command to complete) before exiting, and won't lose messages if restarted.

内容的提问来源于stack exchange,提问作者Hammerbot

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 03:38:18