In distributed systems, microservices frequently need to communicate without tight HTTP coupling. When a customer completes a checkout, inventory must be reserved, payment must be processed, receipt emails must be dispatched, and analytics must be notified.
If you wire all of this through synchronous REST calls, a single timeout cascades into an aborted transaction.
Asynchronous event-driven architecture decouples services by emitting events to a broker. But picking the wrong messaging broker or ignoring dual-write failure modes leads to lost messages and double-charged credit cards.
A few weeks ago, Amir and I redesigned an order processing pipeline for a fintech client. They were trying to use RabbitMQ as an analytical stream while simultaneously struggling with dual-write bugs between their PostgreSQL database and message queues.
Here is the exact decision matrix and implementation blueprint we used to separate Kafka from RabbitMQ and guarantee zero message loss with the Transactional Outbox pattern.
When to Use Kafka vs RabbitMQ
One of the most persistent misconceptions in backend engineering is treating Kafka and RabbitMQ as interchangeable queues. They solve fundamentally different problems:
- RabbitMQ (Smart Broker, Dumb Consumer): An AMQP message broker built for complex task routing, work queues, and priority handling. Once a consumer acknowledges a message, RabbitMQ removes it from memory. Choose RabbitMQ for point-to-point task distribution, background worker jobs, and dead-letter retries.
- Apache Kafka (Dumb Broker, Smart Consumer): An append-only distributed commit log designed for high-throughput stream processing and event sourcing. Kafka retains messages even after they are read, allowing new consumer groups to replay historical streams from any offset. Choose Kafka for analytics pipelines, event sourcing, and high-volume clickstream logs.
Standardizing on CloudEvents 1.0
Never publish unstructured JSON payloads. Without a standardized envelope, consumers break whenever schema versions shift.
We wrap every message in the CloudEvents 1.0 specification:
// schemas/cloudevent.ts
export interface CloudEvent<T = unknown> {
specversion: '1.0';
id: string; // Unique event UUID for consumer deduplication
source: string; // e.g., '/services/orders'
type: string; // e.g., 'com.app.order.created.v1'
time: string; // ISO 8601 timestamp
datacontenttype: 'application/json';
data: T; // Strongly-typed payload
}Preventing Data Inconsistency: The Transactional Outbox Pattern
The classic dual-write bug happens when your service updates the database and then attempts to publish an event:
[Service] ββ 1. UPDATE DB (Success)
β
ββββ 2. PUBLISH Event (Network Timeout / Crash) βIf the database commit succeeds but the message broker disconnects before receiving the event, downstream services never know the order was created. Your database and message queue are now out of sync.
The Transactional Outbox pattern solves this by writing the event directly into an outbox table within the same database transaction:
-- Outbox Table inside PostgreSQL
CREATE TABLE outbox_events (
id UUID PRIMARY KEY,
aggregate_type TEXT NOT NULL,
aggregate_id TEXT NOT NULL,
event_type TEXT NOT NULL,
payload JSONB NOT NULL,
status TEXT NOT NULL DEFAULT 'PENDING' CHECK (status IN ('PENDING', 'PUBLISHED', 'FAILED')),
retry_count INTEGER NOT NULL DEFAULT 0,
created_at TIMESTAMPTZ DEFAULT now() NOT NULL,
published_at TIMESTAMPTZ
);
CREATE INDEX idx_outbox_pending ON outbox_events(status, created_at) WHERE status = 'PENDING';When an order is created, the order record and the outbox event commit in a single ACID transaction:
// services/order-service.ts
import { db } from '../lib/db';
import { randomUUID } from 'crypto';
export async function createOrder(orderData: { userId: string; amount: number }) {
const orderId = randomUUID();
const eventId = randomUUID();
await db.transaction(async (tx) => {
// Write 1: Insert business entity
await tx.query(
'INSERT INTO orders (id, user_id, amount, status) VALUES ($1, $2, $3, $4)',
[orderId, orderData.userId, orderData.amount, 'PENDING']
);
// Write 2: Insert outbox record in the SAME transaction
await tx.query(
`INSERT INTO outbox_events (id, aggregate_type, aggregate_id, event_type, payload)
VALUES ($1, $2, $3, $4, $5)`,
[
eventId,
'order',
orderId,
'com.app.order.created.v1',
JSON.stringify({ orderId, userId: orderData.userId, amount: orderData.amount }),
]
);
});
}A background worker (or Change Data Capture tool like Debezium) polls the pending outbox records, publishes them to Kafka or RabbitMQ, and updates the status to PUBLISHED. If the worker crashes, it retries safely without losing events.
Consumer Idempotency: Handling At-Least-Once Delivery
Distributed brokers guarantee at-least-once delivery, not exactly-once delivery. Network hiccups or consumer crashes before acknowledgment can cause duplicate deliveries.
Every consumer must track processed event IDs before executing side effects:
// consumers/inventory-consumer.ts
import { db } from '../lib/db';
import { CloudEvent } from '../schemas/cloudevent';
export async function processOrderEvent(event: CloudEvent<{ orderId: string }>) {
// Use database unique constraint to guarantee idempotent execution
const processed = await db.query(
'INSERT INTO processed_events (event_id, processed_at) VALUES ($1, NOW()) ON CONFLICT (event_id) DO NOTHING RETURNING event_id;',
[event.id]
);
if (processed.rowCount === 0) {
console.log(`Duplicate event ${event.id} detected. Skipping.`);
return;
}
// Safe to execute side effects
await reserveInventory(event.data.orderId);
}Summary Checklist for Event-Driven Systems
- Choose the right tool: RabbitMQ for routing and task processing; Kafka for replaying streaming data.
- Standardize payloads: Always enforce CloudEvents 1.0 schemas.
- Eliminate dual writes: Write to an outbox table in the same transaction as your state changes.
- Enforce consumer idempotency: Track event IDs to prevent duplicate processing.
Comments
Comments are reviewed before appearing publicly.
No comments yet β be the first.