amqp
@webda/amqp
RabbitMQ-backed Queue and Pub/Sub services for Webda — drop-in replacements for in-memory transports when you need durable, cross-process messaging.
When to use it
- You need a durable task queue backed by RabbitMQ (or any AMQP 0-9-1 broker).
- You need fan-out pub/sub across multiple Webda instances via AMQP exchanges.
- You are replacing
MemoryQueue/MemoryPubSubin a production deployment.
Install
pnpm add @webda/amqp
Configuration
AMQPQueue
{
"services": {
"taskQueue": {
"type": "AMQPQueue",
"url": "amqp://localhost:5672",
"queue": "my-tasks"
}
}
}
| Parameter | Type | Default | Required | Description |
|---|---|---|---|---|
url | string | — | Yes | AMQP broker connection URL (e.g. amqp://user:pass@host:5672) |
queue | string | — | Yes | Name of the AMQP queue to assert and consume |
queueOptions | object | — | No | Options forwarded to channel.assertQueue() (e.g. { durable: true }) |
AMQPPubSubService
{
"services": {
"eventBus": {
"type": "AMQPPubSub",
"url": "amqp://localhost:5672",
"channel": "my-events",
"exchange": { "type": "fanout", "durable": true }
}
}
}
| Parameter | Type | Default | Required | Description |
|---|---|---|---|---|
url | string | — | Yes | AMQP broker connection URL |
channel | string | — | Yes | Exchange name to publish/subscribe to |
subscription | string | "" | No | Subscription queue name (auto-generated exclusive queue if empty) |
exchange.type | string | "fanout" | No | Exchange type (fanout, direct, topic, headers) |
exchange.durable | boolean | true | No | Survive broker restarts |
exchange.autoDelete | boolean | false | No | Delete when last binding is removed |
Usage
import { Queue } from "@webda/core";
import { Service } from "@webda/core";
import { Bean, Inject } from "@webda/core";
@Bean
export class OrderService extends Service {
@Inject("taskQueue")
queue: Queue<{ orderId: string }>;
async placeOrder(orderId: string): Promise<void> {
// Publish a task to RabbitMQ
await this.queue.sendMessage({ orderId });
}
}
// Consumer side: start a queue worker via AsyncJobService or
// by calling queue.consume(async (msg) => { ... })
Reference
- API reference: see the auto-generated typedoc at
docs/pages/Modules/amqp/. - Source:
packages/amqp - Related:
@webda/corefor theQueueandPubSubServicebase classes,@webda/asyncfor job orchestration on top of a queue.