AMQP Integration Guide
Use MaestroHub's AMQP connector to exchange messages with any broker that speaks AMQP 1.0 (ISO/IEC 19464) — the open, vendor-neutral standard for enterprise messaging. One connector covers Apache ActiveMQ Artemis, Apache Qpid, Azure Service Bus, Azure Event Hubs, Solace PubSub+, IBM MQ, and RabbitMQ 4.x.
Despite the shared name, AMQP 1.0 and AMQP 0-9-1 are entirely different, incompatible wire protocols. This connector speaks AMQP 1.0. For RabbitMQ 3.x — whose native protocol is AMQP 0-9-1 with the exchange/binding model — use the RabbitMQ connector instead. RabbitMQ 4.x supports AMQP 1.0 natively and works with this connector.
Overview
The AMQP connector delivers:
- Broker-neutral messaging — one connector for every standard AMQP 1.0 broker
- At-least-once settlement — accepted messages are removed; failed dispatches are released for redelivery or dead-lettered by the broker
- Durable queue consumption with credit-based flow control (prefetch)
- Request/Reply over auto-created dynamic reply queues with correlation-ID matching
- One-shot receive and non-destructive browse for queue draining and inspection, returning as soon as the queue stops delivering (no fixed wait) with optional SQS-style long polling
- SASL ANONYMOUS / PLAIN / EXTERNAL authentication and TLS/mTLS with custom certificates
- AMQP over WebSocket (
ws:///wss://) for brokers behind HTTP-only ingress - RabbitMQ virtual host selection for non-default vhosts
- Payload templating with dynamic
((parameter))syntax
Connection Configuration
Creating an AMQP Connection
Navigate to Connections → New Connection → AMQP and fill in these details.
AMQP Connection Creation Fields
1. Profile Information
| Field | Default | Description |
|---|---|---|
| Profile Name | - | A descriptive name for this connection profile (required, max 100 characters) |
| Description | - | Optional description for this AMQP connection |
2. Broker (Connection tab)
| Field | Type | Default | Required | Description |
|---|---|---|---|---|
| Endpoint URL | String | amqp://localhost:5672 | Yes | Broker endpoint. amqp:// (plaintext, port 5672), amqps:// (TLS, port 5671), or ws:// / wss:// for AMQP-over-WebSocket endpoints (Azure Service Bus, brokers behind HTTP ingress). |
| Container ID | String | maestrohub | No | AMQP container identifier reported to the broker; shows up in broker-side connection listings |
| Virtual Host (RabbitMQ) | String | - | No | RabbitMQ 4.x only — selects the virtual host to connect to. Leave empty for RabbitMQ's default vhost (/), and always leave empty for other brokers, which use the endpoint hostname. |
| Connect Timeout (seconds) | Number | 30 | No | Maximum time for the TCP/TLS dial and AMQP open handshake (1–300) |
3. Authentication (Security tab)
| Field | Type | Default | Description |
|---|---|---|---|
| SASL Mechanism | Select | anonymous | anonymous (no credentials — test/dev brokers), plain (username/password), or external (identity from the TLS client certificate — mTLS) |
| Username | String | - | Required when SASL Mechanism is plain |
| Password | Password | - | Stored encrypted. Leave empty on edit to keep the existing value. |
4. TLS / mTLS (Security tab)
| Field | Type | Default | Description |
|---|---|---|---|
| Verify TLS Certificate | Boolean | false | Verify the broker's TLS certificate (recommended for production). Off by default — industrial brokers commonly run self-signed certificates. When enabled, the SNI name and custom CA fields become available. |
| TLS Server Name (SNI) | String | - | Server name for SNI / certificate verification (defaults to the endpoint host). Shown only while verification is enabled. |
| CA Certificate (PEM) | String | - | Custom CA to verify the broker's server certificate (self-signed brokers). Defaults to the system trust store. Shown only while verification is enabled. |
| Client Certificate (PEM) | Password | - | Optional client certificate for mutual TLS. Required for SASL external. Always available — mTLS authenticates you to the broker and works independently of server-certificate verification. |
| Client Key (PEM) | Password | - | Optional client private key for mutual TLS. Must be provided together with the certificate. |
5. Connection Labels
| Field | Default | Description |
|---|---|---|
| Labels | - | Key-value pairs to categorize and organize this AMQP connection (max 10 labels) |
- TLS is derived from the endpoint scheme:
amqps://andwss://endpoints negotiate TLS,amqp://andws://do not. There is no separate TLS on/off switch. - SASL
externalrequires a TLS endpoint (amqps:///wss://) plus a client certificate and key — the certificate is the identity. - AMQP URLs can carry credentials (
amqp://user:pass@host); prefer the SASL fields so the password is stored encrypted. - The default endpoint works for a local broker during development. For Azure Service Bus use
amqps://<namespace>.servicebus.windows.netor thewss://WebSocket form.
RabbitMQ 4.x enforces its v2 address grammar on every function's Address field: use /queues/<queue-name> for queues and /exchanges/<exchange>/<routing-key> to publish through an exchange. A bare queue name is refused with amqp:invalid-field (amqp_address_v1_not_permitted). RabbitMQ also does not auto-declare queues over AMQP 1.0 — create them first (management UI or API). Artemis and most other brokers accept plain queue/topic names and auto-create on first use.
Function Builder
Creating AMQP Functions
After the connection is configured:
- Open the connection and go to its Functions tab → New Function
- Choose Publish, Publish Batch, Request, Receive, Browse, or Subscribe as the function type
- Configure the address and operation-specific fields
Publish Function
Purpose: Send a single message to a queue or topic and wait for the broker's disposition — success means the broker settled the transfer, not just that bytes left the socket. Publish is Store & Forward eligible: with durability enabled on the pipeline, messages are buffered locally during a broker outage and drained when it returns.
Configuration Fields
| Field | Type | Required | Default | Description |
|---|---|---|---|---|
| Address | String | Yes | - | Queue or topic address exactly as the broker names it. Supports ((parameter)) syntax. |
| Payload | String | Yes | - | Message payload. Supports ((parameter)) syntax. Empty payloads are rejected before reaching the broker. |
| Durable | Boolean | No | true | Mark the message durable so the broker persists it before acknowledging |
| Content Type | String | No | - | MIME content type stamped on the message (e.g. application/json). Supports ((parameter)) syntax. |
| Subject | String | No | - | AMQP message subject property — a routing hint for some brokers. Supports ((parameter)) syntax. |
| Correlation ID | String | No | - | Correlation identifier for request/response tracking. Supports ((parameter)) syntax. |
| Application Properties (JSON) | Object | No | - | Custom key/value properties carried in the message's application-properties section. Supports ((parameter)) syntax. |
Use Cases: Publish line events to an enterprise service bus, send commands to a device-facing queue, forward pipeline output to an ERP integration queue
Publish Batch Function
Purpose: Send an array of messages over a single sender link in one operation, amortizing link-attach overhead for bulk emission. Each element becomes one AMQP message and is individually settled; sending stops at the first broker rejection and the result reports how many were sent. Also Store & Forward eligible.
Configuration Fields
| Field | Type | Required | Default | Description |
|---|---|---|---|---|
| Address | String | Yes | - | Target address. Supports ((parameter)) syntax. |
| Messages (JSON array) | Array | Yes | - | Array of message bodies — strings or objects, e.g. [{"n": 1}, {"n": 2}, "plain text"]. Supports ((parameter)) syntax. |
| Durable | Boolean | No | true | Mark every message in the batch durable |
| Content Type | String | No | - | MIME content type stamped on every message |
| Subject | String | No | - | AMQP message subject property stamped on every message in the batch. Supports ((parameter)) syntax. |
| Correlation ID | String | No | - | Correlation identifier stamped on every message in the batch. Supports ((parameter)) syntax. |
| Application Properties (JSON) | Object | No | - | Custom key/value properties carried in every message's application-properties section. Supports ((parameter)) syntax. |
Use Cases: Bulk-forward buffered telemetry, ship a batch of transformed records downstream
Request Function
Purpose: Send a request and block until the correlated reply arrives on an auto-created dynamic reply queue, or the timeout elapses. Implements the standard AMQP request/reply pattern — the responder must copy the request's message ID into the reply's correlation ID.
Configuration Fields
| Field | Type | Required | Default | Validation | Description |
|---|---|---|---|---|---|
| Address | String | Yes | - | - | Address the responder listens on. Supports ((parameter)) syntax. |
| Payload | String | Yes | - | - | Request payload. Supports ((parameter)) syntax. |
| Content Type | String | No | - | - | MIME content type stamped on the request |
| Timeout (ms) | Number | No | 1800000 | 1000–3600000 | Maximum time to wait for the correlated reply |
Use Cases: RPC-style calls to a service listening on a queue, synchronous command/response with an ERP adapter
Request needs broker support for dynamic (server-named) reply queues, which RabbitMQ added in 4.1. On older 4.0 servers the request fails with Dynamic source not supported. Artemis, Service Bus, and the other brokers support dynamic termini out of the box.
Receive Function
Purpose: Pull up to N messages from a queue in one shot, accepting each one (removing it from the queue). Returns as soon as the queue stops delivering — an empty result on an empty queue is near-immediate success, not an error. Use for scheduled queue-drain pipelines where a continuous subscription is not wanted.
Configuration Fields
| Field | Type | Required | Default | Validation | Description |
|---|---|---|---|---|---|
| Address | String | Yes | - | - | Source address. Supports ((parameter)) syntax. |
| Max Messages | Number | No | 10 | 1–10000 | Upper bound on messages pulled in this call — a cap, not a target; the call never waits to "fill up" |
| Wait Time (seconds) | Number | No | 0 | 0–300 | Long-polling: block up to this many seconds for the first message — covers both an empty queue and a slow/remote broker. 0 returns immediately with whatever the queue holds. |
| Idle Cutoff (ms) | Number | No | 500 | 100–10000 | Advanced: a silent gap this long ends the pull ("queue is dry"). Raise it for slow or high-latency brokers that may pause mid-delivery; lowering it speeds up returns but risks partial batches. |
Use Cases: Drain a batch queue every five minutes from a scheduled pipeline, fetch pending commands for sequential processing
Browse Function
Purpose: Peek up to N messages without consuming them — every message is released back to the broker, so it remains on the queue for real consumers. Same Max Messages / Wait Time / Idle Cutoff fields as Receive.
Use Cases: Inspect a dead-letter queue without draining it, verify a producer is publishing before wiring a consumer
While a browse is in flight the peeked messages are temporarily held by the browsing receiver and invisible to competing consumers; their redelivery order afterwards is broker-dependent. This is inherent to the emulation — the AMQP 1.0 non-destructive-read terminus is not exposed by the underlying client library.
Subscribe Function
Purpose: Stream every message on a queue or topic into the pipeline as a trigger event. Credit-based prefetch bounds how many unsettled messages are in flight; each message is accepted only after the pipeline has durably captured it — failures release it for redelivery or dead-letter it, giving at-least-once delivery. Subscribe functions drive pipeline triggers — they cannot be invoked manually.
Configuration Fields
| Field | Type | Required | Default | Validation | Description |
|---|---|---|---|---|---|
| Address | String | Yes | - | - | Source address exactly as the broker names it |
| Prefetch | Number | No | 10 | 1–1000 | Receiver-link credit — how many unsettled messages to hold in flight at once |
Use Cases: Consume ERP/MES events from a Service Bus queue, trigger a pipeline from an Artemis work queue, load-balance consumption across replicas via competing consumers
Using Parameters
AMQP publish, batch, request, receive, and browse functions support parameterized values via the ((parameterName)) syntax in their templatable fields (marked with (()) in the form). Parameters are auto-detected — typing ((orderId)) in the payload automatically adds orderId to the function's parameter list, which then appears as an input on the pipeline node.
The connector templating engine recognises ((double parentheses)) — {{ curly braces }} are ignored and stored as literal text.
Pipeline Integration
Use the AMQP functions you create here as nodes inside the Pipeline Designer. Drag in the AMQP Publish, Publish Batch, Request, Receive, or Browse nodes to act on brokers from inside a pipeline, or the AMQP Trigger node to start pipelines automatically on incoming messages.
Common Use Cases
ERP/MES Event Bridge
Subscribe to a Service Bus or Artemis queue carrying ERP business events, transform them with downstream nodes, and publish results into the Unified Namespace. Competing consumers on the same queue load-balance across MaestroHub replicas automatically.
Guaranteed Delivery to the Enterprise Bus
Publish pipeline output to an enterprise AMQP queue with Store & Forward durability: during a broker outage messages buffer locally and drain automatically when the broker returns — no data loss, no manual replay.
Scheduled Queue Drain
A schedule-triggered pipeline calls Receive every few minutes to pull pending work items in bounded batches, instead of holding a permanent subscription. An empty queue simply completes with zero messages.
Dead-Letter Inspection
A monitoring pipeline Browses the dead-letter queue and raises a notification when its count grows — without consuming the dead-lettered messages, so operators can still inspect and reprocess them.
Synchronous Command over the Bus
Use Request to send a command to a device- or service-facing queue and wait for the correlated reply, with a timeout bounded per call.
Troubleshooting
"Attach refused: amqp_address_v1_not_permitted" (RabbitMQ)
The Address field contains a bare queue name. RabbitMQ 4.x only accepts its v2 grammar: /queues/<queue-name> or /exchanges/<exchange>/<routing-key>. Update the address on the function. The connection self-heals automatically after this error — no restart needed.
"no queue '…' in vhost" (RabbitMQ)
The queue does not exist. RabbitMQ does not auto-declare queues over AMQP 1.0 — create it via the management UI/API first. If the queue exists in a non-default virtual host, set Virtual Host (RabbitMQ) on the connection.
"Dynamic source not supported" on Request
The broker cannot create the dynamic reply queue. On RabbitMQ this requires version 4.1 or newer; on other brokers check that server-named/temporary queues are permitted for your user.
"SASL PLAIN auth failed"
The username/password pair was rejected. Note RabbitMQ's default guest user is only allowed to connect from localhost. Re-enter the password on the connection (a masked value left untouched keeps the previously stored secret).
Receive/Browse returns fewer messages than expected
The pull ends when the broker stops delivering for the Idle Cutoff window. On slow or high-latency links raise Wait Time (seconds) (covers the first message) and/or Idle Cutoff (ms) (covers gaps mid-stream). Remember Max Messages is a cap — the call does not wait for the queue to reach it.