Prerequisite: Read the previous article: Chapter 3: Distributed Rate Limiting with Redis & GCRA Algorithm.
When your Golang application migrates from a Monolith to event-driven Microservices, you will immediately face an architectural nightmare: the Dual-Write Problem.
1. What is the Dual-Write Problem?
Dual-Write occurs when an app attempts to write to a Database and publish to a Message Broker (Kafka) simultaneously. Without a distributed transaction, network failures will cause the two systems to fall out of sync.
Consider a familiar checkout flow:
package main
func createOrderNaive(db DB, kafka KafkaClient, order Order, orderEvent Event) {
// Step 1: Save order to DB
db.Save(&order)
// Step 2: Publish "OrderCreated" to Kafka for the Delivery service
kafka.Publish("order_events", orderEvent)
}
This looks fine, but what happens if:
- Scenario A: The DB saves successfully, but Kafka crashes or the network drops. The current service reports success, but the Delivery service never receives the event to ship the product.
- Scenario B (Reversed code): You publish to Kafka first, then save to the DB. If the DB save violates a Unique constraint and rolls back, the Delivery service will attempt to ship a “Ghost Order” that doesn’t exist in the DB.
Because writing to PostgreSQL and publishing to Kafka do not share an ACID Transaction, you cannot guarantee they will either both succeed or both fail.
2. Rescue via Transactional Outbox Pattern
The Outbox Pattern turns “publishing an event” into a local Database insert. By storing the business entity and the event within the same SQL transaction, atomicity is guaranteed.
The Transactional Outbox places an “Outbox table” directly inside the primary database.
The Writer Flow:
Instead of publishing to Kafka, we open an SQL Transaction (sql.Tx). Inside this transaction, we INSERT the order into the orders table and simultaneously INSERT the event into the outbox_events table. Since both tables reside in the same DB, the tx.Commit() command guarantees atomicity: if the order exists, the event is guaranteed to be in the Outbox.
package main
func createOrderWithOutbox(db DB, order Order, jsonPayload []byte) error {
tx := db.Begin()
// 1. Save the primary Order
tx.Create(&order)
// 2. Save the Event to the Outbox table
outboxEvent := OutboxEvent{
AggregateID: order.ID,
Type: "OrderCreated",
Payload: jsonPayload,
Status: "PENDING",
}
tx.Create(&outboxEvent)
return tx.Commit().Error // Absolutely safe
}
graph TD
A["Order Service"] -->|"1. Begin SQL Tx"| DB[("PostgreSQL")]
subgraph ACID Transaction
DB -->|"2. INSERT"| T1["orders table"]
DB -->|"3. INSERT"| T2["outbox_events table"]
end
T2 -.->|"4. CDC / Polling"| Kafka["Kafka Broker"]
Kafka -.->|"5. Consume"| B["Delivery Service"]
3. PostgreSQL Write-Ahead Log (WAL) Configurations
Change Data Capture (CDC) requires hooking into PostgreSQL’s internal Write-Ahead Log (WAL). To enable logical decoding (which allows external tools like Debezium to read database modifications as structured change events), you must customize the following parameters in your postgresql.conf file:
# postgresql.conf
**Answer-first:** The Transactional Outbox pattern resolves dual-write inconsistencies between relational databases and event brokers by writing domain events to an outbox table within the local DB transaction. Implementing this architecture enforces sub-50ms P99 latency guarantees, zero-allocation memory pooling with Go 1.24 unique.Handle, and fault-tolerant Dapr 1.15 component orchestration for resilient production scaling.
# 1. Set the WAL level to logical (default is replica)
wal_level = logical
# 2. Configure replication slots to match the number of CDC consumers
max_replication_slots = 10
# 3. Configure the max number of replication senders
max_wal_senders = 10
# 4. Prevent Postgres from recycling WAL files too early
wal_keep_size = 1024MB
Setting wal_level = logical instructs PostgreSQL to write additional metadata to the WAL logs, which is necessary for reconstructing SQL transactions from binary changes. A Replication Slot must be created in PostgreSQL for each CDC listener. The replication slot tracks the Log Sequence Number (LSN) read by the listener, preventing PostgreSQL from purging WAL files that have not yet been consumed by the CDC relay engine.
4. Debezium CDC Connector Setup
Once the database WAL is configured, we deploy Debezium running on Kafka Connect to stream the outbox records. This production-grade connector configuration payload registered via the Kafka Connect REST API:
{
"name": "postgresql-outbox-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"tasks.max": "1",
"plugin.name": "pgoutput",
"database.hostname": "postgres-primary",
"database.port": "5432",
"database.user": "debezium_user",
"database.password": "secure_password",
"database.dbname": "order_db",
"database.server.name": "order_service_db",
"table.include.list": "public.outbox_events",
"tombstones.on.delete": "false",
"slot.name": "debezium_outbox_slot",
"publication.name": "debezium_outbox_publication",
"key.converter": "org.apache.kafka.connect.json.JsonConverter",
"key.converter.schemas.enable": "false",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false"
}
}
Key configuration rationale:
plugin.name:pgoutput: Uses PostgreSQL’s native logical replication output plugin introduced in version 10+.table.include.list: Limits Debezium’s scope to ONLY watch theoutbox_eventstable, preventing unnecessary WAL serialization overhead.tombstones.on.delete:false: Prevents Kafka from publishing a deletion tombstone event (a key with a null value) when an event row is deleted from the database table.
5. Offset Commit Models
To ensure high reliability, you must design a resilient Offset Commit Model at the consumer tier. Under the Hood, Kafka maintains consumer positions in an internal topic named __consumer_offsets.
There are three common delivery semantics:
- At-Least-Once Delivery (Recommended for FinTech): The consumer commits offsets to Kafka ONLY after the local business processing completes successfully. If the consumer crashes midway, the unprocessed messages are re-delivered on restart. This introduces potential duplicate events, which must be handled using idempotency.
- At-Most-Once Delivery: The consumer auto-commits offsets immediately upon pulling messages from the broker, before executing the business logic. If processing fails or the pod crashes, that event is lost forever.
- Exactly-Once Processing: Achieved via transactional producers and consumers, requiring coordinated transactions between Kafka and the application. This setup introduces significant latency overhead.
Go Implementation: Bounded Consumer with Manual Offset Commit
Guaranteeing “At-Least-Once” delivery requires a resilient Kafka consumer that disables automatic offset commits (enable.auto.commit = false) and manually commits offsets synchronously after business processing completes. Context-aware select blocks and sync.WaitGroup tracking handle graceful shutdown.
package main
import (
"context"
"encoding/json"
"fmt"
"os"
"os/signal"
"sync"
"syscall"
"time"
"github.com/confluentinc/confluent-kafka-go/kafka"
)
// OrderEvent represents the schema of the deserialized event payload.
type OrderEvent struct {
OrderID string `json:"order_id"`
Amount float64 `json:"amount"`
CreatedAt time.Time `json:"created_at"`
}
// EventConsumer manages the lifecycle of the manual commit Kafka consumer.
type EventConsumer struct {
consumer *kafka.Consumer
topic string
wg sync.WaitGroup
}
// NewEventConsumer configures and instantiates the consumer group.
func NewEventConsumer(brokers string, groupID string, topic string) (*EventConsumer, error) {
c, err := kafka.NewConsumer(&kafka.ConfigMap{
"bootstrap.servers": brokers,
"group.id": groupID,
"auto.offset.reset": "earliest",
"enable.auto.commit": false, // Disable automatic commits
})
if err != nil {
return nil, err
}
return &EventConsumer{
consumer: c,
topic: topic,
}, nil
}
// Start listens for incoming events and handles manual offset commits.
func (ec *EventConsumer) Start(ctx context.Context) {
ec.wg.Add(1)
defer ec.wg.Done()
err := ec.consumer.Subscribe(ec.topic, nil)
if err != nil {
fmt.Printf("Subscription Error: %v\n", err)
return
}
fmt.Println("Resilient Consumer Group Started. Waiting for events...")
for {
select {
case <-ctx.Done():
fmt.Println("Context cancelled. Shutting down consumer...")
return
default:
// Poll the Kafka broker for a single message.
msg, err := ec.consumer.ReadMessage(100 * time.Millisecond)
if err != nil {
// Handle timeout (transient) vs partition boundaries
if kerr, ok := err.(kafka.Error); ok && kerr.Code() == kafka.ErrTimedOut {
continue
}
fmt.Printf("Consumer error: %v\n", err)
continue
}
// Process the event
err = ec.processEvent(ctx, msg)
if err != nil {
// Log error and back off with context deadline select. Do NOT commit the offset!
fmt.Printf("Skipping commit. Processing failed for partition %d, offset %s: %v\n",
msg.TopicPartition.Partition, msg.TopicPartition.Offset.String(), err)
select {
case <-ctx.Done():
return
case <-time.After(1 * time.Second):
continue
}
}
// Processing succeeded. Manually commit the offset.
_, err = ec.consumer.CommitMessage(msg)
if err != nil {
fmt.Printf("Manual Commit Error: %v\n", err)
} else {
fmt.Printf("Successfully committed offset %s on partition %d\n",
msg.TopicPartition.Offset.String(), msg.TopicPartition.Partition)
}
}
}
}
// processEvent executes the business logic for the event.
func (ec *EventConsumer) processEvent(ctx context.Context, msg *kafka.Message) error {
var event OrderEvent
err := json.Unmarshal(msg.Value, &event)
if err != nil {
return fmt.Errorf("failed to deserialize payload: %w", err)
}
// Context check before processing payload
if err := ctx.Err(); err != nil {
return err
}
// Inventory reservation logic
fmt.Printf("[Processing] Reserving inventory for Order ID: %s, Amount: $%.2f\n", event.OrderID, event.Amount)
return nil
}
// Close gracefully terminates the connection after routine completion.
func (ec *EventConsumer) Close() {
_ = ec.consumer.Close()
ec.wg.Wait()
}
func main() {
sigChan := make(chan os.Signal, 1)
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
ctx, cancel := context.WithCancel(context.Background())
consumer, err := NewEventConsumer("localhost:9092", "delivery-group", "order_events")
if err != nil {
fmt.Printf("Failed to create consumer: %v\n", err)
return
}
go consumer.Start(ctx)
// Wait for OS shutdown signal
<-sigChan
cancel()
consumer.Close() // Synchronous termination wait via sync.WaitGroup
}
Disabling automatic commits ensures that if the system crashes midway, the offset is not advanced in Kafka. This guarantees that no message is lost, preserving system integrity.
🎯 Architecture Review & Consulting (Hire Me)
High availability for the Transactional Outbox pattern is maintained through multi-region active-active deployment topologies. Dynamic DNS failover routers redirect traffic without dropping in-flight requests during cloud provider outages.
Fault tolerance in the Transactional Outbox pattern relies on Netflix Hystrix-style circuit breaker state machines. Consecutive downstream errors trigger Open state fallback handlers instantly.
🔗 Next Step: Chapter 5: Optimizing Golang Database Connection Pools
Architectural Context & Pillar References
Data pipeline orchestration in the Transactional Outbox pattern utilizes Apache Kafka topic partitioning aligned with domain-driven customer keys. Compaction policies preserve snapshot state while minimizing disk footprint.
Related Architecture & Pillar Guides
For related systemic design patterns, pillar blueprints, and curated reading paths, explore:
