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: truemeans 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
phpcommand incomplete.
Key Changes Needed
- Enable manual message acknowledgments: Replace
noAck: truewithnoAck: false, and only acknowledge messages after your child process completes successfully. This ensures RabbitMQ requeues the message if your process crashes before finishing. - Track shutdown state: Add a flag to signal when the process should stop accepting new messages.
- Listen for shutdown signals: Catch
SIGINT(Ctrl+C) andSIGTERM(system stop commands) to trigger graceful shutdown. - 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:
isShuttingDownprevents new messages from being processed once a shutdown signal is received. Any messages received after shutdown starts are rejected and requeued. - Signal Handling: When
SIGINTorSIGTERMis 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
closeResourcesfunction 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
相关产品推荐
相关产品推荐

