Outbox
An outbox stores the integration events of a change in the same transaction as the change itself. A relay then reads them and hands them to an event publisher, retrying until it succeeds: an event is never lost, even if the broker is down when the change is saved.
await this.outbox.add(order.pullDomainEvents().map((event) => this.translator.translate(event, { correlationId })));
await new OutboxRelay(outbox, publisher).relay();When to use
Whenever an event must not be lost: other bounded contexts, other systems or an audit trail depend on it. Add the events in the unit of work of the change, and run the relay in the background.
Usage
Add the events of a change
In the command handler, inside the unit of work, after saving the aggregate, translated by an event translator.
return this.unitOfWork.run(async () => {
await this.orders.save(order);
await this.outbox.add(order.pullDomainEvents().map((event) => this.translator.translate(event, { correlationId: orderId })));
return ok();
});Relay the events
OutboxRelay reads a batch of pending events, publishes them and marks them as published. It returns how many it published; if publishing fails, it throws and the events stay pending for the next call. Build it in the composition root with a batch size, and call it on a schedule from a driving adapter: a timer, a cron job, a worker.
const relay = new OutboxRelay(outbox, publisher, 100);import type { OutboxRelay } from "@alveolus/core";
export class OutboxRelayJob {
constructor(private readonly relay: OutboxRelay) {}
start(): NodeJS.Timeout {
return setInterval(() => this.tick(), 1000);
}
private async tick(): Promise<void> {
try {
await this.relay.relay();
} catch (error) {
console.error("The outbox relay failed; the events stay pending for the next tick.", error);
}
}
}Implement it
Extend Outbox in a driven adapter, typically over a table written through the transaction of the unit of work. Integration events are JSON: store them as they are, for instance in a jsonb column, and return them unchanged from pending, oldest first.
import { type AnyIntegrationEvent, Outbox } from "@alveolus/core";
import type { Pool } from "pg";
export class PgOutbox extends Outbox {
constructor(private readonly db: Pool) {
super();
}
async add(events: readonly AnyIntegrationEvent[]): Promise<void> {
for (const event of events) {
await this.db.query("INSERT INTO outbox (id, event) VALUES ($1, $2)", [event.id, event]);
}
}
async pending(limit: number): Promise<readonly AnyIntegrationEvent[]> {
const { rows } = await this.db.query("SELECT event FROM outbox WHERE published_at IS NULL ORDER BY created_at LIMIT $1", [limit]);
return rows.map((row) => row.event);
}
async markPublished(ids: readonly string[]): Promise<void> {
await this.db.query("UPDATE outbox SET published_at = now() WHERE id = ANY($1)", [ids]);
}
}Do not publish directly
await this.orders.save(order);
await this.publisher.publish(events);await this.unitOfWork.run(async () => {
await this.orders.save(order);
await this.outbox.add(events);
return ok();
});Why?
If the process stops between save and publish, the order is placed and nobody hears about it. If the broker is down, the command fails although the order is valid. Writing the events in the same transaction as the change, then relaying them, removes both cases.
Reference
abstract class Outbox extends Port {
abstract add(events: readonly AnyIntegrationEvent[]): Promise<void>;
abstract pending(limit: number): Promise<readonly AnyIntegrationEvent[]>;
abstract markPublished(ids: readonly string[]): Promise<void>;
}
class OutboxRelay {
constructor(outbox: Outbox, publisher: EventPublisher, batchSize?: number);
relay(): Promise<number>;
}| Member | Type | Description |
|---|---|---|
add(events) | Promise<void> | Stores the events, in the current transaction. |
pending(limit) | Promise<readonly AnyIntegrationEvent[]> | Returns up to limit unpublished events, oldest first. |
markPublished(ids) | Promise<void> | Marks the events as published. |
new OutboxRelay(outbox, publisher, batchSize = 100) | Relays from the outbox to the publisher. | |
relay() | Promise<number> | Publishes one batch and returns how many events it published; 0 when none is pending. |
Caveats
- Delivery is at least once: if the relay stops between publishing and marking, the events are published again. Consumers ignore duplicates by
id. - Several relays running at once may publish the same batch; lock the rows in
pending(for instanceFOR UPDATE SKIP LOCKED) if you run more than one.
Import from @alveolus/core or @alveolus/core/outbox.
See also
- Unit of Work, the transaction the outbox is written in
- Event translators, which build what the outbox stores
- Event publishers, what the relay calls