
Introduction
When each microservice has its own database, the biggest challenge is ensuring data consistency between services. This lesson delves into patterns that address dual write, reliable event publishing, and eventual consistency.
1. Dual Write Problem
1.1 Problem
When a service needs to both update database and publish event, if either fails → inconsistency:
Scenario 1: DB thành công, Event thất bại
┌──────────────┐
│ Order Service│
│ │
│ 1. Save to DB ✅ (Order created)
│ 2. Publish to Kafka ❌ (Network error)
│ │
│ → Order tồn tại trong DB nhưng không ai biết
└──────────────┘
Scenario 2: Event thành công, DB thất bại
┌──────────────┐
│ Order Service│
│ │
│ 1. Publish to Kafka ✅ (OrderCreated sent)
│ 2. Save to DB ❌ (DB connection failed)
│ │
│ → Downstream services xử lý order không tồn tại
└──────────────┘
1.2 Why not use distributed transactions?
// ❌ Anti-pattern: Distributed transaction giữa DB và Kafka
@Transactional
public void createOrder(Order order) {
database.save(order); // Transaction participant 1
kafka.publish(orderCreated); // Transaction participant 2
// Kafka không support XA/2PC → KHÔNG THỂ atomic
}
2. Outbox Pattern
2.1 Concepts
Record the event to the outbox table in the same database transaction as the business data. A separate process (Relay/Poller) reads the outbox and publishes it to the message broker:
┌─────────────────────────────────────────────────┐
│ Order Service │
│ │
│ ┌───────────────────────────────────────────┐ │
│ │ Database Transaction (ACID) │ │
│ │ │ │
│ │ INSERT INTO orders (...) VALUES (...); │ │
│ │ INSERT INTO outbox (...) VALUES (...); │ │
│ │ │ │
│ │ → Both succeed or both fail │ │
│ └───────────────────────────────────────────┘ │
│ │
│ ┌───────────────────────────────────────────┐ │
│ │ Outbox Relay (separate process) │ │
│ │ │ │
│ │ 1. Poll outbox table │ │
│ │ 2. Publish event to Kafka │ │
│ │ 3. Mark outbox entry as published │ │
│ └────────────────────┬──────────────────────┘ │
└───────────────────────┼──────────────────────────┘
│
▼
┌─────────┐
│ Kafka │
└─────────┘
2.2 Implementation
Outbox Table:
CREATE TABLE outbox (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
aggregate_type VARCHAR(255) NOT NULL, -- "Order"
aggregate_id VARCHAR(255) NOT NULL, -- "O-001"
event_type VARCHAR(255) NOT NULL, -- "OrderCreated"
payload JSONB NOT NULL, -- event data
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
published_at TIMESTAMPTZ, -- NULL = chưa publish
retries INT DEFAULT 0
);
CREATE INDEX idx_outbox_unpublished
ON outbox (created_at) WHERE published_at IS NULL;
Business Logic:
@Service
public class OrderService {
@Transactional
public Order createOrder(CreateOrderCommand cmd) {
// 1. Save business data
Order order = new Order(cmd);
order.setStatus(OrderStatus.PENDING);
orderRepository.save(order);
// 2. Save outbox event (same transaction!)
OutboxEvent event = OutboxEvent.builder()
.aggregateType("Order")
.aggregateId(order.getId())
.eventType("OrderCreated")
.payload(toJson(new OrderCreatedPayload(order)))
.build();
outboxRepository.save(event);
return order;
}
}
Outbox Relay (Polling):
@Scheduled(fixedDelay = 500) // Poll mỗi 500ms
public void publishOutboxEvents() {
List<OutboxEvent> events = outboxRepository
.findUnpublishedOrderByCreatedAt(BATCH_SIZE);
for (OutboxEvent event : events) {
try {
kafkaTemplate.send(
event.getAggregateType().toLowerCase() + ".events",
event.getAggregateId(),
event.getPayload()
).get(); // Wait for Kafka ACK
event.setPublishedAt(Instant.now());
outboxRepository.save(event);
} catch (Exception e) {
event.incrementRetries();
if (event.getRetries() >= MAX_RETRIES) {
event.moveToDeadLetter();
}
outboxRepository.save(event);
}
}
}
2.3 Advantages and disadvantages
| Advantages | Disadvantages |
|---|---|
| Guaranteed delivery (at-least-once) | Polling delay (latency) |
| Atomic with business data | Outbox table needs cleanup |
| Simple implementation | At-least-once → consumer must be idempotent |
| No need to change database technology | Add storage overhead |
3. Change Data Capture (CDC)
3.1 Concepts
CDC captures database changes (INSERT, UPDATE, DELETE) from transaction log (WAL in PostgreSQL) and streams directly to the message broker:
┌──────────────────────────────────────────────┐
│ Order Service │
│ │
│ ┌────────────┐ ┌──────────────────────┐ │
│ │ Application│───▶│ PostgreSQL │ │
│ │ │ │ ┌──────────────────┐ │ │
│ └────────────┘ │ │ orders table │ │ │
│ │ └──────────────────┘ │ │
│ │ ┌──────────────────┐ │ │
│ │ │ outbox table │ │ │
│ │ └──────────────────┘ │ │
│ │ ┌──────────────────┐ │ │
│ │ │ WAL (Write-Ahead │ │ │
│ │ │ Log) │ │ │
│ │ └────────┬─────────┘ │ │
│ └──────────┼───────────┘ │
└───────────────────────────────┼───────────────┘
│
┌───────────▼───────────┐
│ Debezium Connector │
│ (reads WAL stream) │
└───────────┬───────────┘
│
┌───────────▼───────────┐
│ Apache Kafka │
│ topic: outbox.events │
└───────────────────────┘
3.2 Debezium Configuration
{
"name": "order-outbox-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "order-db",
"database.port": "5432",
"database.user": "debezium",
"database.password": "${file:/secrets/db-password}",
"database.dbname": "orders",
"database.server.name": "order-service",
"plugin.name": "pgoutput",
"table.include.list": "public.outbox",
"transforms": "outbox",
"transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
"transforms.outbox.table.field.event.id": "id",
"transforms.outbox.table.field.event.key": "aggregate_id",
"transforms.outbox.table.field.event.type": "event_type",
"transforms.outbox.table.field.event.payload": "payload",
"transforms.outbox.route.topic.replacement": "${routedByValue}.events",
"transforms.outbox.table.expand.json.payload": true
}
}
3.3 CDC vs Polling
| Criteria | Polling (Outbox Relay) | CDC (Debezium) |
|---|---|---|
| Latency | 100ms-1s (poll interval) | ~10ms (near real-time) |
| Database load | Frequent queries | Reads WAL (minimal impact) |
| Infrastructure | Simple (no extra components) | Need Debezium + Kafka Connect |
| Complexity | Low | Medium |
| Reliability | Good | Excellent (WAL = source of truth) |
| Recommended | Small/Medium scale | Large scale, low latency |
4. Idempotent Consumers
4.1 Why do we need Idempotency?
At-least-once delivery can send duplicate messages:
Producer ──msg──▶ Broker ──msg──▶ Consumer
│
Process ✅
│
ACK ──▶ Broker (network timeout!)
│
Broker re-delivers msg
│
Consumer receives msg AGAIN
Process AGAIN? → Duplicate!
4.2 Implementation Strategies
Strategy 1: Idempotency Key Table
CREATE TABLE processed_events (
event_id UUID PRIMARY KEY,
processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
-- Consumer logic
BEGIN;
-- Check if already processed
INSERT INTO processed_events (event_id) VALUES ($event_id)
ON CONFLICT (event_id) DO NOTHING;
-- If inserted (not duplicate), process the event
IF FOUND THEN
-- Business logic here
UPDATE inventory SET quantity = quantity - $qty WHERE product_id = $pid;
END IF;
COMMIT;
Strategy 2: Natural Idempotency
// SET operations are naturally idempotent
// "Set order status to CONFIRMED" — same result regardless of how many times
@EventHandler
public void on(PaymentCompletedEvent event) {
Order order = orderRepository.findById(event.getOrderId());
// Idempotent check: if already confirmed, skip
if (order.getStatus() == OrderStatus.CONFIRMED) {
return;
}
order.setStatus(OrderStatus.CONFIRMED);
orderRepository.save(order);
}
Strategy 3: Version-based (Optimistic Locking)
UPDATE orders
SET status = 'CONFIRMED', version = version + 1
WHERE id = 'O-001' AND version = 3;
-- If version doesn't match → someone else already updated → skip
-- Rows affected = 0 → duplicate processing detected
5. Eventual Consistency Strategies
5.1 Read Your Own Writes
Make sure users always see the data they just changed:
Approach 1: Sticky Session
→ Route user requests đến cùng replica đã nhận write
Approach 2: Client-side Optimistic Update
→ UI cập nhật ngay trước khi server confirm, rollback nếu fail
Approach 3: Write-through Read Model
→ Sau khi write, đồng bộ invalidate read cache
Client ──POST /orders──▶ Order Service ──▶ Write DB
│
Client ──GET /orders/123──▶ Order Service │
↑ │ │
│ Check: event synced? │
│ Yes ──▶ Read DB │
│ No ──▶ Write DB│ (fallback)
└────────────────────────────────────────────┘
5.2 Consistency Monitoring
-- Monitor replication lag
SELECT
topic,
partition,
MAX(offset) - MIN(committed_offset) AS lag
FROM consumer_offsets
GROUP BY topic, partition
HAVING lag > 1000; -- Alert if lag > 1000 messages
# Prometheus alert for consistency lag
- alert: HighConsumerLag
expr: kafka_consumer_group_lag > 5000
for: 5m
labels:
severity: warning
annotations:
summary: "Consumer lag is high"
description: "Consumer group {{ $labels.group }} has lag {{ $value }}"
6. End-to-End Pattern
6.1 Recommended Architecture
┌─────────────────────────────────────────────────────────────┐
│ End-to-End Data Flow │
│ │
│ ┌──────────┐ ┌──────────────────────┐ ┌───────────┐ │
│ │ Service │ │ PostgreSQL │ │ Debezium │ │
│ │ Logic │───▶│ ┌────────┐ ┌───────┐ │───▶│ CDC │ │
│ │ │ │ │Business│ │Outbox │ │ │ Connector │ │
│ └──────────┘ │ │ Table │ │ Table │ │ └─────┬─────┘ │
│ │ └────────┘ └───────┘ │ │ │
│ └──────────────────────┘ │ │
│ ▼ │
│ ┌──────────┐ │
│ │ Kafka │ │
│ │ Topic │ │
│ └────┬─────┘ │
│ │ │
│ ┌────────────────────┼────┐ │
│ │ │ │ │
│ ┌─────▼────┐ ┌─────▼────┐ │
│ │ Consumer │ │ Consumer │ │
│ │ Idempot. │ │ Idempot. │ │
│ │ Check ✓ │ │ Check ✓ │ │
│ └──────────┘ └──────────┘ │
└─────────────────────────────────────────────────────────────┘
Guarantees:
✅ Atomic write (business data + outbox in same TX)
✅ At-least-once delivery (CDC from WAL)
✅ Exactly-once processing (idempotent consumers)
✅ Near real-time (<50ms latency)
Summary
- Dual Write Problem is the root cause of data inconsistency in microservices
- Outbox Pattern is solved by writing events to the same DB transaction
- CDC (Debezium) reads WAL stream for near real-time event delivery
- Idempotent consumers ensures exact-once processing despite at-least-once delivery
- Eventual consistency is unavoidable — design for it instead of against it
- Monitoring consistency lag is required for production