Skip to main content

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 / MemoryPubSub in a production deployment.

Install​

pnpm add @webda/amqp

Configuration​

AMQPQueue​

{
"services": {
"taskQueue": {
"type": "AMQPQueue",
"url": "amqp://localhost:5672",
"queue": "my-tasks"
}
}
}
ParameterTypeDefaultRequiredDescription
urlstring—YesAMQP broker connection URL (e.g. amqp://user:pass@host:5672)
queuestring—YesName of the AMQP queue to assert and consume
queueOptionsobject—NoOptions 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 }
}
}
}
ParameterTypeDefaultRequiredDescription
urlstring—YesAMQP broker connection URL
channelstring—YesExchange name to publish/subscribe to
subscriptionstring""NoSubscription queue name (auto-generated exclusive queue if empty)
exchange.typestring"fanout"NoExchange type (fanout, direct, topic, headers)
exchange.durablebooleantrueNoSurvive broker restarts
exchange.autoDeletebooleanfalseNoDelete 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/core for the Queue and PubSubService base classes, @webda/async for job orchestration on top of a queue.

Classes​