@inteli.city/node-red-contrib-rabbit-mq 2.0.1
Node-RED nodes for producing, consuming and acknowledging RabbitMQ messages
@inteli.city/node-red-contrib-rabbit-mq
Node-RED nodes for consuming, producing, and acknowledging RabbitMQ messages via AMQP.
Table of Contents
- Overview
- Nodes
- Basic Flow
- System Behavior
- Important Decisions
- Common Issues
- Debug
- Testing
- Simple Example
Overview
These three nodes implement the consume → process → ack pattern for RabbitMQ.
mq.consumerconnects to a queue and emits messages into the Node-RED flowmq.ackwacknowledges the message after processing — requiredmq.producerpublishes messages to an exchange with a routing key
The consumer operates in manual ack mode. RabbitMQ only removes a message from the queue when mq.ackw is called. Without an ack, the message stays in-flight until the timeout fires and is then requeued.
Nodes
mq.consumer
Connects to a RabbitMQ queue and emits one message per delivery.
Parameters
| Field | Description |
|---|---|
| Host | Broker address |
| User / Password | AMQP credentials |
| Queue | Queue name (asserted as durable on connect) |
| Prefetch | Maximum number of unacknowledged messages at once |
| Timeout (min) | Maximum time before automatic nack. 0 uses the default of 120 minutes (2 hours). Any positive value is used as-is |
| Heartbeat (s) | AMQP heartbeat, detects dead TCP links. Default 30; 0 leaves it to the broker |
| Health check (s) | Consumer liveness watchdog interval. Default 30; 0 disables it |
| SSL/TLS | Use amqps:// — required for cloud brokers |
| Debug mode | Enables verbose logs |
Behavior
- Connects automatically on deploy
- On failure: reconnects with exponential backoff (1s → 2s → ... → 60s)
- Prefetch limits parallelism: with
prefetch=1, the consumer only receives the next message after acking the current one - Timeout: if the flow does not call
mq.ackwwithin the configured window, the message is automatically nacked and requeued
Lifecycle status
The status is green only when all three conditions hold: the AMQP connection is open, the channel is open, and the consumer subscription is registered. Any other state is visible in the status text, so a stalled consumer is never shown as healthy:
| Status | Meaning |
|---|---|
connecting… (#n) |
Opening the TCP/AMQP connection (attempt n) |
connection up, opening channel… |
AMQP connection established, channel not ready |
channel open, registering consumer… |
Queue asserted and prefetch set, subscription not registered yet |
consuming: <queue> (green) |
Connection + channel + consumer all live |
<reason> – retry in n s (red) |
Recovery pending; see the reason table below |
Failure reasons: tcp-refused, dns-failure, tcp-lost, heartbeat-timeout, handshake-timeout, amqp-connection-error, amqp-connection-closed, channel-error, channel-closed, channel-unavailable, consumer-cancelled, consumer-not-registered, auth-refused, queue-config-conflict, queue-not-found, queue-locked.
auth-refused and queue-config-conflict require manual intervention: they are logged with node.error and retried at the maximum 60 s interval, so fixing the broker side recovers the node without a redeploy.
Recovery guarantees
- A channel that dies while the TCP connection stays open triggers full recovery
- Broker-side consumer cancellation (
basic.cancel, e.g. queue deleted) triggers full recovery - The watchdog detects a channel or subscription that became unusable without emitting any event
- Repeated
error/closeevents collapse into a single reconnect — no duplicate connections, channels, consumers or timers - Inflight entries and their timers belonging to a dead channel are discarded (RabbitMQ requeues those messages)
Output
Each emitted message contains:
msg.payload → message content (JSON or string)
msg.rabbitmq → internal context (do not modify)
msg._msgid → delivery identifier
Warning:
msg.rabbitmqmust be passed intact tomq.ackw. Any node that recreates themsgobject (e.g.return { payload: ... }) will destroy this context.
mq.ackw
Acknowledges a message to RabbitMQ.
When to use
Always, after processing a message consumed by mq.consumer.
What it does
- Calls
channel.ack(message)on the broker - Cancels the message timeout timer
- Passes
msgdownstream (nodes can be chained after the ack)
Consequences of not using it
- The message remains "unacked" in RabbitMQ
- The prefetch window fills up after N messages (where N = configured prefetch)
- The consumer silently stops receiving new messages
- After the timeout, the message is nacked and requeued — the cycle repeats
mq.producer
Publishes messages to a RabbitMQ exchange.
Parameters
| Field | Description |
|---|---|
| Host | Broker address |
| User / Password | AMQP credentials |
| Exchange | Exchange name (asserted as direct and durable on connect) |
| Routing Key | Default routing key |
| SSL/TLS | Use amqps:// — required for cloud brokers |
| Debug mode | Enables publish logs |
Connection behavior
- Connects on deploy. On failure: exponential backoff (1s → 30s)
- Uses a confirm channel: only advances in the flow after the broker confirms receipt
msg.exchangeandmsg.routingKeyoverride the configured values per message
When disconnected
Messages are dropped. If the producer is not connected when a message arrives, the message is lost and a
warnis written to the logs.
There is no internal buffer. Messages arriving during reconnection are discarded.
Basic Flow
mq.consumer → [processing] → mq.ackw
mq.consumerreceives the message and emits it- The flow processes it (function node, HTTP request, database, etc.)
mq.ackwconfirms to RabbitMQ that the message was handled
Ack is mandatory. Without it:
message delivered → not acked → timeout (≥5 min) → nack → requeued → delivered again
System Behavior
Automatic reconnect All nodes reconnect automatically with exponential backoff. Restarting Node-RED is not required.
Prefetch
Controls how many messages can be in-flight simultaneously per consumer. With prefetch=1, the flow is strictly sequential. With prefetch=N, up to N messages can be processed in parallel — but all must be acked.
Blocking due to missing ack If all prefetch slots are occupied by unacked messages, RabbitMQ stops delivering new messages to that consumer. No error is raised — the consumer simply goes silent.
Redelivery
Nacked or timed-out messages are requeued and redelivered. The field msg.rabbitmq.message.fields.redelivered indicates whether a message is a redelivery.
Important Decisions
The producer drops messages when disconnected. There is no internal buffer. If guaranteed delivery is required, implement persistence in the flow before the producer.
The system is at-least-once. The same message may be delivered more than once (after nack, reconnection, or timeout). Flows should be idempotent or detect duplicates.
Ack is the flow's responsibility.
mq.consumerdoes not ack automatically. Every path through the flow must reachmq.ackw— including error paths.
Common Issues
Consumer stopped receiving messages while the status is green → Prefetch is saturated. Since the status only stays green while the connection, channel and subscription are all live, a green node that receives nothing means the prefetch window is full. Most common cause:
mq.ackwis not being reached in the flow (disconnected node, silent error, ormsg.rabbitmqdestroyed by a function node). Enable debug mode and checkinflight sizeagainst the configured prefetch. → In versions before 2.1.0 this could also mean the channel had died or the broker had cancelled the consumer: both were only logged and left the node green forever, requiring a redeploy.Consumer keeps cycling through
retry in …→ Read the reason in the status text:tcp-refused/dns-failureare network or hostname problems,auth-refusedis credentials,queue-config-conflictmeans the queue exists with different arguments thandurable: true.Duplicate messages → Expected behavior. Indicates a nack or reconnection occurred before the ack.
Messages disappearing in the producer → The producer was disconnected. Messages arriving during reconnection are discarded.
msg.rabbitmqis undefined in mq.ackw → An upstream node recreated themsgobject without preservingmsg.rabbitmq. Fix by usingmsg.payload = ...; return msg;instead ofreturn { payload: ... }.Consumer connects but receives no messages after restart → Check if the queue has messages with "unacked" status from a previous session. They will be requeued once the old connection expires on the broker.
Debug
Debug mode is enabled via checkbox in mq.consumer and mq.producer.
When enabled, the following logs appear:
mq.consumer:
- Connection details (URL, channel, queue, consumerTag)
- Every lifecycle transition (
state=… – …), including which failure caused each teardown - Health check results (
healthcheck ok – queue=… messages=… consumers=… inflight=…) - Discarded stale attempts, collapsed duplicate events and skipped reconnects
- Broker flow control (
connection blocked/unblocked) - Each message received (msgid, deliveryTag)
- Inflight size and prefetch state
- 10s timer to detect flows missing an ack
Warnings and errors are always logged, regardless of debug mode: every teardown reports its reason ([consumer] channel-closed – tearing down and scheduling reconnect), and discarded inflight messages are counted.
mq.producer:
- Connection established
- Each publish (exchange, routing key)
When to enable:
- To investigate why the consumer stopped receiving messages
- To confirm that messages are being published
- To trace the lifecycle of a specific message
Disable in production. Generates one log entry per message, creating excessive noise and unnecessary I/O.
Testing
npm test
Runs the consumer lifecycle suite (Node's built-in test runner, no extra dependencies) against a fake amqplib with mocked timers: connection loss, channel death with the connection still open, broker cancellation, watchdog detection, event storms, redeploy races, backoff, and the unchanged ack/auto-nack behaviour.
test/VALIDATION.md documents the matching manual scenarios against a real broker (broker restart, network interruption, queue deletion, incompatible queue configuration, leak checks).
Simple Example
Publishing messages
[inject (repeat: 1s)] → [mq.producer]
Configure mq.producer:
- Exchange:
X.Teste - Routing Key:
R.Teste - SSL: enabled
Consuming and acknowledging messages
[mq.consumer] → [debug] → [mq.ackw]
Configure mq.consumer:
- Queue:
X.Teste - Prefetch:
1 - Timeout:
0(uses the 2-hour default) - SSL: enabled
The debug node displays msg.payload. The mq.ackw confirms to the broker. Without mq.ackw at the end, the consumer will stall after the first message (prefetch=1).