#golang#microservices#architecture#transactional-outbox#postgresql#distributed-systems

Transactional Outbox Pattern di Go: Solusi Masalah Dual-Write

"Panduan arsitektur Transactional Outbox Pattern di Go untuk mengatasi masalah Dual-Write pada distributed systems dan microservices."

By huud
0 views
~5 min read
Transactional Outbox Pattern di Go: Solusi Masalah Dual-Write

Transactional Outbox Pattern di Go: Solusi Masalah Dual-Write

Dalam arsitektur microservices yang mengadopsi Event-Driven Architecture, kebutuhan umum yang sering muncul adalah: memperbarui database lokal DAN mengirim pesan/event ke message broker (Kafka, NATS, RabbitMQ) dalam satu alur bisnis.

Contoh kasus pada pembuatan order:

  1. Simpan order baru ke PostgreSQL dengan status CREATED.

  2. Kirim event OrderCreated ke Kafka agar Service Inventory dapat memotong stok.

Apa yang terjadi jika langkah 1 sukses namun pengiriman pesan ke Kafka di langkah 2 gagal karena network error? Database lokal menyimpan order, tapi inventory tidak pernah tahu (data inconsistency). Sebaliknya, jika Anda mengirim event ke Kafka terlebih dahulu, lalu database crash saat menyimpan order, inventory memotong stok untuk order yang tidak pernah ada. Masalah ini dikenal sebagai The Dual-Write Problem.

Solusi standar industri untuk masalah ini adalah Transactional Outbox Pattern.


1. Konsep Transactional Outbox

Alih-alih mengirim event langsung ke broker, kita menyimpan event tersebut ke dalam tabel database lokal (outbox) di dalam transaksi database ACID yang sama dengan data bisnis utama.

  1. Buka database transaction (BEGIN).

  2. Insert data ke tabel orders.

  3. Insert event payload ke tabel outbox_events.

  4. Commit database transaction (COMMIT).

Kedua operasi dijamin atomik: keduanya sukses bersama atau rollback bersama.

Selanjutnya, sebuah background worker terpisah (Relay Worker) membaca record dari tabel outbox_events dan mempublikasikannya ke message broker. Setelah broker mengonfirmasi penerimaan pesan (ACK), worker menandai record outbox sebagai PUBLISHED atau menghapusnya.


2. Database Schema

sql
CREATE TABLE orders (
    id UUID PRIMARY KEY,
    customer_id VARCHAR(64) NOT NULL,
    total_amount NUMERIC(12, 2) NOT NULL,
    status VARCHAR(32) NOT NULL,
    created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
);

CREATE TABLE outbox_events (
    id UUID PRIMARY KEY,
    aggregate_type VARCHAR(64) NOT NULL,
    aggregate_id VARCHAR(64) NOT NULL,
    event_type VARCHAR(64) NOT NULL,
    payload JSONB NOT NULL,
    status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
    retry_count INT NOT NULL DEFAULT 0,
    created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
    published_at TIMESTAMP WITH TIME ZONE NULL
);

CREATE INDEX idx_outbox_pending ON outbox_events(status, created_at) WHERE status = 'PENDING';

3. Implementasi Transaksi Atomik di Go

Berikut implementasi repository pattern yang menyatukan insert order dan event outbox:

go
package order

import (
	"context"
	"database/sql"
	"encoding/json"
	"fmt"
	"time"

	"github.com/google/uuid"
)

type Order struct {
	ID          string  `json:"id"`
	CustomerID  string  `json:"customer_id"`
	TotalAmount float64 `json:"total_amount"`
	Status      string  `json:"status"`
}

type Service struct {
	db *sql.DB
}

func NewService(db *sql.DB) *Service {
	return &Service{db: db}
}

func (s *Service) CreateOrder(ctx context.Context, customerID string, amount float64) (*Order, error) {
	tx, err := s.db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelReadCommitted})
	if err != nil {
		return nil, fmt.Errorf("gagal memulai db transaction: %w", err)
	}
	defer tx.Rollback()

	order := &Order{
		ID:          uuid.New().String(),
		CustomerID:  customerID,
		TotalAmount: amount,
		Status:      "CREATED",
	}

	// 1. Simpan order ke database
	queryOrder := `INSERT INTO orders (id, customer_id, total_amount, status, created_at) VALUES ($1, $2, $3, $4, $5)`
	_, err = tx.ExecContext(ctx, queryOrder, order.ID, order.CustomerID, order.TotalAmount, order.Status, time.Now().UTC())
	if err != nil {
		return nil, fmt.Errorf("gagal insert order: %w", err)
	}

	// 2. Simpan event ke outbox table dalam transaksi yang sama
	payload, err := json.Marshal(order)
	if err != nil {
		return nil, fmt.Errorf("gagal marshal event payload: %w", err)
	}

	queryOutbox := `
		INSERT INTO outbox_events (id, aggregate_type, aggregate_id, event_type, payload, status, created_at)
		VALUES ($1, $2, $3, $4, $5, 'PENDING', $6)
	`
	eventID := uuid.New().String()
	_, err = tx.ExecContext(ctx, queryOutbox, eventID, "Order", order.ID, "OrderCreated", payload, time.Now().UTC())
	if err != nil {
		return nil, fmt.Errorf("gagal insert outbox event: %w", err)
	}

	// Commit kedua operasi bersamaan
	if err := tx.Commit(); err != nil {
		return nil, fmt.Errorf("commit gagal: %w", err)
	}

	return order, nil
}

4. Background Relay Worker dengan FOR UPDATE SKIP LOCKED

Untuk membaca pesan outbox secara konkuren tanpa race condition antar instance worker, kita gunakan klausa PostgreSQL SELECT ... FOR UPDATE SKIP LOCKED:

go
package outbox

import (
	"context"
	"database/sql"
	"encoding/json"
	"log"
	"time"
)

type Publisher interface {
	Publish(ctx context.Context, topic string, key string, data []byte) error
}

type RelayWorker struct {
	db        *sql.DB
	publisher Publisher
	batchSize int
}

func NewRelayWorker(db *sql.DB, publisher Publisher, batchSize int) *RelayWorker {
	return &RelayWorker{db: db, publisher: publisher, batchSize: batchSize}
}

type EventRecord struct {
	ID        string
	Topic     string
	Key       string
	Payload   []byte
}

func (w *RelayWorker) Start(ctx context.Context, interval time.Duration) {
	ticker := time.NewTicker(interval)
	defer ticker.Stop()

	for {
		select {
		case <-ctx.Done():
			log.Println("Outbox Relay Worker berhenti.")
			return
		case <-ticker.StopChan():
			w.processBatch(ctx)
		}
	}
}

func (w *RelayWorker) processBatch(ctx context.Context) {
	tx, err := w.db.BeginTx(ctx, nil)
	if err != nil {
		log.Printf("Worker begin tx error: %v", err)
		return
	}
	defer tx.Rollback()

	// Kunci row agar tidak diambil oleh replica worker lain
	query := `
		SELECT id, event_type, aggregate_id, payload
		FROM outbox_events
		WHERE status = 'PENDING'
		ORDER BY created_at ASC
		LIMIT $1
		FOR UPDATE SKIP LOCKED
	`

	rows, err := tx.QueryContext(ctx, query, w.batchSize)
	if err != nil {
		log.Printf("Query outbox gagal: %v", err)
		return
	}
	defer rows.Close()

	var events []EventRecord
	for rows.Next() {
		var ev EventRecord
		if err := rows.Scan(&ev.ID, &ev.Topic, &ev.Key, &ev.Payload); err != nil {
			log.Printf("Scan row gagal: %v", err)
			continue
		}
		events = append(events, ev)
	}

	for _, ev := range events {
		// Kirim ke message broker (Kafka/NATS)
		if err := w.publisher.Publish(ctx, ev.Topic, ev.Key, ev.Payload); err != nil {
			log.Printf("Publish event ID %s gagal: %v", ev.ID, err)
			_, _ = tx.ExecContext(ctx, `UPDATE outbox_events SET retry_count = retry_count + 1 WHERE id = $1`, ev.ID)
			continue
		}

		// Tandai sukses
		updateQuery := `UPDATE outbox_events SET status = 'PUBLISHED', published_at = $1 WHERE id = $2`
		_, _ = tx.ExecContext(ctx, updateQuery, time.Now().UTC(), ev.ID)
	}

	_ = tx.Commit()
}

5. Ringkasan Praktik Terbaik

  1. Gunakan SKIP LOCKED: Mencegah deadlock dan lock contention antar multiple pod / replica.

  2. At-Least-Once Delivery: Pola Outbox menjamin pengiriman at-least-once. Pastikan consumer downstream memiliki mekanisme idempotency handler.

  3. Pembersihan Outbox (Data Archiving): Event berstatus PUBLISHED harus di-purge atau dipindahkan ke tabel archive secara berkala (misal tiap 7 hari) agar tabel outbox tetap ramping dan indeks efisien.

  4. Opsi CDC (Change Data Capture): Untuk throughput sangat masif (> 50.000 events/detik), pertimbangkan menggunakan engine CDC seperti Debezium yang membaca WAL PostgreSQL langsung tanpa polling SQL.

About the Author

huud

huud

@huud

About →

Systems architect and software engineer building high-performance distributed platforms.