Frank Fontcha.
← All posts
Qashio Expense Tracker5 min read

Domain events that never get lost: a transactional outbox fanned out to BullMQ

Saving a row and publishing an event are two writes, and one of them will eventually fail. Here's the outbox → relay → one-job-per-handler pipeline I built so budget alerts fire exactly once, even through crashes and retries.

NestJSPostgreSQLBullMQEvent-drivenReliability

In the Qashio expense tracker, creating a transaction has side effects. The budget for that category has to be re-evaluated, the user may need an "80% of budget used" alert, and an activity-log entry gets written. The naive version looks like this:

await transactionsRepo.save(tx);
await eventBus.publish(new TransactionCreated(tx)); // 💥 what if this throws?

That's the dual-write problem. If the process dies, Redis blips or the deploy restarts between those two lines, the transaction exists but nothing downstream ever hears about it. Swap the order and it gets worse: you can alert on a transaction that was rolled back.

The fix is old and boring and it works: a transactional outbox. The event goes into the same database transaction as the data, and a separate relay moves it to the queue afterwards.

The flow end to end

  1. 1Use caseCreateTransactionUseCase opens a unit of work and inserts the transaction row.
  2. 2PublisherIn the same Postgres transaction, it inserts an outbox_events row: type transaction.created plus a JSON payload.
  3. 3PostgresCOMMIT. Data and event become durable together, or neither does.
  4. 4Outbox relayEvery second: SELECT … FOR UPDATE SKIP LOCKED LIMIT 100 on unpublished rows.
  5. 5Outbox relayBuilds one BullMQ job per (event × handler), each with a deterministic jobId, and calls addBulk. Then it marks the rows published_at in the same DB transaction.
  6. 6BullMQ workerRuns exactly one handler per job. A throw is retried with backoff, up to 5 attempts.
  7. 7Budgets listenerRecomputes usage and emits budget.threshold_reached only on an upward crossing of 80% or 100%.
  8. 8Outbox → NotificationsThat event goes through the same outbox, and the notifications module sends the in-app and email alert.

1. Write the event with the data

The trick that makes this clean in NestJS is @nestjs-cls/transactional. The repository and the outbox publisher both resolve the current transaction from async-local storage, so the use case never threads an EntityManager through every call:

create-transaction.use-case.ts (trimmed)
return this.unitOfWork.run(async () => {
  const tx = Transaction.create(input);          // domain entity, Decimal amounts
  await this.transactions.save(tx);
  await this.events.publish(                      // writes to outbox_events,
    new TransactionCreated(tx.snapshot()),        // not to Redis
  );
  return tx;
});

If anything inside run() throws, both rows roll back. Here, "publishing" means inserting a row.

2. Relay with SKIP LOCKED

The relay is a polling loop. Polling sounds crude, but a 1-second tick over an indexed published_at IS NULL column costs almost nothing, and it has no moving parts to break.

outbox-relay.ts
return this.dataSource.transaction(async (manager) => {
  const rows = await manager
    .createQueryBuilder(OutboxEventOrmEntity, "o")
    .where("o.published_at IS NULL")
    .orderBy("o.created_at", "ASC")
    .limit(OUTBOX_BATCH_SIZE)
    .setLock("pessimistic_write")
    .setOnLocked("skip_locked")        // other API replicas take other rows
    .getMany();
  if (rows.length === 0) return 0;
 
  const jobs = toHandlerJobs(rows, this.handlers);   // one job per (row, handler)
  await Promise.race([this.queue.addBulk(jobs), timeout(ENQUEUE_TIMEOUT_MS)]);
 
  await manager.update(OutboxEventOrmEntity,
    { id: In(rows.map((r) => r.id)) }, { publishedAt: new Date() });
  return rows.length;
});

Three details carry the weight here:

  • FOR UPDATE SKIP LOCKED lets you run N API instances without a leader election. Each replica grabs a disjoint batch, and nobody blocks.
  • The enqueue has a timeout race. If Redis hangs, the DB transaction rolls back, the rows stay unpublished and the next tick retries them. Nothing is lost; it's only late.
  • Marking published_at happens after addBulk succeeds, in the same transaction that holds the lock. A crash between the two means the batch gets enqueued twice, which the next section makes harmless.

3. One job per handler, with deterministic job IDs

Most event-bus tutorials enqueue one job per event and let a dispatcher call every subscriber. That breaks on retries: if the notifications handler fails, retrying the job re-runs the budget handler too.

So the relay fans out. At startup, a registry uses Nest's DiscoveryService to find every method decorated with @OnDomainEvent("transaction.created") and gives each one a stable ID like TransactionEventsListener.onCreated. Each outbox row becomes one job per handler:

const jobId = domainEventJobId(row.id, handler.id); // e.g. "evt:9f2c…:TransactionEventsListener.onCreated"

BullMQ ignores an add whose jobId already exists. So when a crashed relay re-enqueues a batch, the duplicate jobs are dropped at the queue. Each handler retries on its own, five attempts with backoff, and a handler that no longer exists throws BullMQ's UnrecoverableError instead of burning retries. When a job runs out of attempts, it's logged at error level and sent to Sentry.

4. Make the consumer idempotent and edge-triggered

The budget listener turns create, update and delete events into { before, after } snapshots and asks one question: did this change cross a threshold, going up?

budget-thresholds.ts
export const BUDGET_ALERT_THRESHOLDS = [100, 80]; // highest first
 
export function crossedThreshold(prev: Decimal, curr: Decimal, limit: Decimal) {
  const before = usedPercent(prev, limit);
  const after = usedPercent(curr, limit);
  return BUDGET_ALERT_THRESHOLDS.find((t) => before.lt(t) && after.gte(t)) ?? null;
}

Edge-triggering matters. A level check like after >= 80 would re-alert on every transaction once a budget is over 80%. Then the emitted budget.threshold_reached event carries a dedupeKey derived from the source event ID, and the insert uses ON CONFLICT DO NOTHING. A retried job produces the same key, so the user gets one alert, not five.

All money math goes through decimal.js. 0.1 + 0.2 has no place in a budget percentage.

What I'd tell you before you build one

  • Start with polling. Debezium/CDC is great when you need it, but a 30-line relay with SKIP LOCKED covers most products and is trivial to debug: the outbox table is your event log.
  • Clean up. Published rows are deleted after 7 days by a scheduled job. Without that, the table grows forever.
  • Fan out per handler. Isolated retries are worth the extra jobs.
  • Test the replays. Write specs that deliver the same event twice and that cross a threshold from both directions, alongside the happy path. Duplicate delivery isn't an edge case in this design; it's the contract.
  • The brief asked for Kafka. The domain didn't need a distributed log, and an outbox on Postgres plus Redis gave the same guarantee with two fewer moving parts. Kafka can replace BullMQ later without touching a single use case, since they only ever talk to the outbox.

Written by Frank Donald Kamga Fontcha

Senior Full Stack Developer · Lead Software Engineer, Dubai, UAE. Questions, or want this pattern in your stack? Email me.