#golang#microservices#nats#jetstream#event-driven#architecture

Event-Driven Architecture di Go dengan NATS Core dan JetStream

"Membangun arsitektur event-driven berkecepatan tinggi di Go menggunakan NATS Core dan JetStream dengan worker pool dan deduplication."

By huud
0 views
~5 min read
Event-Driven Architecture di Go dengan NATS Core dan JetStream

Event-Driven Architecture di Go dengan NATS Core dan JetStream

Dalam arsitektur microservices, synchronous RPC (seperti HTTP REST atau gRPC) seringkali menciptakan tight coupling dan cascading failures. Jika Service A memanggil Service B, dan Service B memanggil Service C, keterlambatan atau kegagalan di Service C akan langsung merambat ke Service A.

Event-Driven Architecture (EDA) menyelesaikan masalah ini dengan membalik arah dependensi: service memublikasikan event saat terjadi perubahan state, dan service lain bereaksi terhadap event tersebut secara asinkron.

NATS adalah message broker berbasis Go berkinerja tinggi, berbobot ringan (single binary, footprint memori < 30MB), dan mendukung at-most-once (NATS Core) hingga at-least-once / exactly-once deduplication (NATS JetStream).


1. NATS Core vs NATS JetStream

  • NATS Core: Pub/Sub murni fire-and-forget dalam memori. Latensi ultra-rendah (< 50 mikrosekon), throughput jutaan pesan/detik. Tidak ada disk persistence. Cocok untuk telemetri, metrics, atau distributed cache invalidation.

  • NATS JetStream: Persistence layer bawaan NATS. Menyimpan stream pesan di NVMe/SSD, mendukung consumer groups (Queue Groups), ACK/NAK semantics, deduplication window, dan replay riwayat pesan. Cocok untuk order processing, payments, dan event sourcing.


2. Inisialisasi Koneksi NATS Idiomatik di Go

Koneksi NATS harus menggunakan reconnection backoff dan context handling yang tepat:

go
package natsbroker

import (
	"fmt"
	"log"
	"time"

	"github.com/nats-io/nats.go"
)

type Client struct {
	Conn *nats.Conn
	JS   nats.JetStreamContext
}

func NewClient(url string) (*Client, error) {
	opts := []nats.Option{
		nats.Name("order-service-worker"),
		nats.ReconnectWait(2 * time.Second),
		nats.MaxReconnects(-1), // Unlimited retry
		nats.DisconnectErrHandler(func(c *nats.Conn, err error) {
			log.Printf("NATS disconnected: %v. Reconnecting...", err)
		}),
		nats.ReconnectHandler(func(c *nats.Conn) {
			log.Printf("NATS reconnected ke: %s", c.ConnectedUrl())
		}),
		nats.ClosedHandler(func(c *nats.Conn) {
			log.Printf("NATS connection closed.")
		}),
	}

	nc, err := nats.Connect(url, opts...)
	if err != nil {
		return nil, fmt.Errorf("koneksi NATS gagal: %w", err)
	}

	js, err := nc.JetStream(nats.PublishAsyncMaxPending(256))
	if err != nil {
		nc.Close()
		return nil, fmt.Errorf("inisialisasi JetStream gagal: %w", err)
	}

	return &Client{Conn: nc, JS: js}, nil
}

3. Publisher: Menerbitkan Event dengan Deduplication

Untuk menjamin idempotent publishing saat retry, kita sematkan Nats-Msg-Id pada header pesan JetStream.

go
package natsbroker

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

	"github.com/nats-io/nats.go"
)

type OrderCreatedEvent struct {
	OrderID   string    `json:"order_id"`
	UserID    string    `json:"user_id"`
	Amount    float64   `json:"amount"`
	Timestamp time.Time `json:"timestamp"`
}

func (c *Client) PublishOrderCreated(ctx context.Context, event OrderCreatedEvent) error {
	data, err := json.Marshal(event)
	if err != nil {
		return fmt.Errorf("marshal payload gagal: %w", err)
	}

	msg := &nats.Msg{
		Subject: "ORDERS.created",
		Data:    data,
		Header:  nats.Header{},
	}
	// Menggunakan OrderID sebagai message ID untuk deduplication JetStream
	msg.Header.Set(nats.MsgIdHdr, fmt.Sprintf("order-created-%s", event.OrderID))

	ack, err := c.JS.PublishMsg(msg, nats.Context(ctx))
	if err != nil {
		return fmt.Errorf("publish message gagal: %w", err)
	}

	log.Printf("Event published ke stream %s sequence %d", ack.Stream, ack.Sequence)
	return nil
}

4. Consumer: Pull-Based Worker Pool dengan Manual ACK

Pull-based consumer memberi kontrol penuh kepada Go service untuk mengatur kecepatan konsumsi data (backpressure control), mencegah memory bloat saat beban lonjakan data (traffic spikes).

go
package natsbroker

import (
	"context"
	"encoding/json"
	"errors"
	"log"
	"time"

	"github.com/nats-io/nats.go"
)

type Consumer struct {
	client *Client
}

func NewConsumer(client *Client) *Consumer {
	return &Consumer{client: client}
}

func (c *Consumer) StartWorkerPool(ctx context.Context, workers int) error {
	// Pastikan Stream sudah ada
	_, err := c.client.JS.AddStream(&nats.StreamConfig{
		Name:      "ORDERS",
		Subjects:  []string{"ORDERS.*"},
		Storage:   nats.FileStorage,
		Retention: nats.LimitsPolicy,
		Duplicates: 5 * time.Minute,
	})
	if err != nil && !errors.Is(err, nats.ErrStreamNameAlreadyInUse) {
		return fmt.Errorf("create stream gagal: %w", err)
	}

	// Pull subscription pada durable consumer
	sub, err := c.client.JS.PullSubscribe(
		"ORDERS.created",
		"payment-processor-group",
		nats.ManualAck(),
		nats.AckWait(10*time.Second),
	)
	if err != nil {
		return fmt.Errorf("pull subscribe gagal: %w", err)
	}

	for i := 0; i < workers; i++ {
		workerID := i + 1
		go c.workerLoop(ctx, workerID, sub)
	}

	return nil
}

func (c *Consumer) workerLoop(ctx context.Context, id int, sub *nats.Subscription) {
	log.Printf("[Worker %d] siap memproses event", id)

	for {
		select {
		case <-ctx.Done():
			log.Printf("[Worker %d] berhenti...", id)
			return
		default:
			// Fetch batch 1 pesan dengan timeout 2 detik
			msgs, err := sub.Fetch(1, nats.MaxWait(2*time.Second))
			if err != nil {
				if errors.Is(err, nats.ErrTimeout) {
					continue
				}
				log.Printf("[Worker %d] fetch error: %v", id, err)
				time.Sleep(500 * time.Millisecond)
				continue
			}

			for _, msg := range msgs {
				c.handleMessage(ctx, id, msg)
			}
		}
	}
}

func (c *Consumer) handleMessage(ctx context.Context, workerID int, msg *nats.Msg) {
	var event OrderCreatedEvent
	if err := json.Unmarshal(msg.Data, &event); err != nil {
		log.Printf("[Worker %d] JSON parsing error: %v. Mengirim Terminate.", workerID, err)
		_ = msg.Term() // Jangan di-retry jika schema rusak
		return
	}

	log.Printf("[Worker %d] Memproses order %s senilai %.2f", workerID, event.OrderID, event.Amount)

	// Simulasi pemrosesan bisnis
	if err := processPayment(ctx, event); err != nil {
		log.Printf("[Worker %d] Gagal memproses payment: %v. Mengirim NAK.", workerID, err)
		_ = msg.NakWithDelay(2 * time.Second) // Retry kembali setelah delay
		return
	}

	// Sukses, kirim acknowledgement
	if err := msg.Ack(); err != nil {
		log.Printf("[Worker %d] ACK gagal: %v", workerID, err)
	}
}

func processPayment(ctx context.Context, event OrderCreatedEvent) error {
	// Business logic di sini
	return nil
}

5. Menangani Poison Pill Message dan Dead Letter Queue (DLQ)

Pesan yang rusak (malformed payload) atau terus-menerus gagal diproses bisnis dapat membuat antrean macet. NATS JetStream menyediakan metadata pengiriman:

go
meta, err := msg.Metadata()
if err == nil && meta.NumDelivered > 5 {
    log.Printf("Pesan %s gagal lebih dari 5 kali. Pindah ke DLQ.", msg.Subject)
    _ = msg.Term() // Hentikan retry
    // Simpan ke DB / table Dead Letter untuk investigasi manual
}

6. Ringkasan Praktik Terbaik

  1. Gunakan Pull Consumer untuk Workers: Push consumer dapat membanjiri memory worker jika load melonjak tinggi.

  2. Set MsgIdHdr: Mencegah duplicate processing jika publisher melakukan retry saat network blip.

  3. Pilih Storage yang Tepat: Gunakan FileStorage untuk data transaksi dan MemoryStorage untuk ephemeral cache stream.

  4. Batasi Retry dengan Term(): Gunakan msg.Term() untuk pesan yang corrupt secara skema agar tidak terjadi infinite loop retry.

About the Author

huud

huud

@huud

About →

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