AizuDemy

Tutorial Build Event-Driven Microservice Pakai NATS JetStream & Go

Tutorial Build Event-Driven Microservice Pakai NATS JetStream & Go
๐ŸŽง
Dengarkan Artikel Ini
Suara AI Otomatis โ€ข 7 mnt baca baca
โšก TL;DR

Poin Kunci Artikel Ini:

  • Event-driven architecture (EDA) menyelesaikan masalah tersebut dengan memisahkan pengirim pesan (producer) dan penerima pesan (consumer) secara asinkron.
  • Jika consumer tidak aktif saat pesan dipublikasikan, pesan tersebut hilang (at-most-once delivery).
  • NATS JetStream menghadirkan engine persistence bawaan pada NATS Server.
๐Ÿ“‹ Daftar Isi Materi Tutup โ–ด

1. Latar Belakang & Arsitektur NATS JetStream

Arsitektur microservice berbasis pemanggilan HTTP/REST synchronous memunculkan masalah tight coupling, cascading failure, dan latensi tinggi. Event-driven architecture (EDA) menyelesaikan masalah tersebut dengan memisahkan pengirim pesan (producer) dan penerima pesan (consumer) secara asinkron. Sistem berbasis event membutuhkan message broker yang tangguh, cepat, dan ringan.

Anatomi NATS Core vs NATS JetStream

NATS Core beroperasi sebagai pub/sub messaging system in-memory tanpa mekanisme persistence. Jika consumer tidak aktif saat pesan dipublikasikan, pesan tersebut hilang (at-most-once delivery). NATS JetStream menghadirkan engine persistence bawaan pada NATS Server. JetStream menambahkan kapabilitas esensial untuk sistem enterprise:

  • Message Persistence: Penyimpanan pesan pada disk atau memory berkinerja tinggi.
  • At-Least-Once Delivery: Garansi pengiriman ulang pesan sampai consumer mengirim konfirmasi (acknowledgement).
  • Durable & Ephemeral Consumer State: Pelacakan posisi konsumsi pesan (sequence offset) per consumer.
  • Deduplication Engine: Pencegahan pemrosesan pesan ganda menggunakan ID unik dalam window waktu tertentu.
  • Raft Consensus: Replikasi data antarnode cluster tanpa tergantung dependensi eksternal.

Perbandingan Mendalam: NATS JetStream vs Apache Kafka vs RabbitMQ

KriteriaNATS JetStreamApache KafkaRabbitMQ
Resource FootprintSangat rendah (Single binary Go, RAM < 50MB)Tinggi (Butuh JVM, RAM > 2GB, KRaft/ZooKeeper)Sedang (Erlang runtime, RAM ~200MB)
Latency (p99)Sub-milidetik (< 1ms)5-10ms (Optimasi batch throughput)2-5ms (Tergantung routing key)
OperasionalTanpa dependensi, auto-cluster via RaftKompleks (Partition balancing, log compaction)Sedang (Erlang cluster cookie, queue mirroring)
Model ConsumerPush & Pull ConsumerPull Consumer (Consumer Group)Push Consumer (Smart Broker, Dumb Consumer)
Ordering GuaranteePer Subject dalam StreamPer Partition dalam TopicPer Queue

2. Setup Server Infrastructure & Stream Configuration

Konfigurasi NATS Server via Docker Compose

Jalankan NATS Server dengan fitur JetStream aktif menggunakan docker-compose.yml berikut. Konfigurasi ini mengaktifkan engine JetStream (-js) dan menentukan direktori penyimpanan data (-sd /data):

version: '3.8'

services:
  nats:
    image: nats:2.10-alpine
    container_name: nats-jetstream
    command: ["-js", "-sd", "/data", "-m", "8222"]
    ports:
      - "4222:4222" # Client Port
      - "8222:8222" # HTTP Management & Metrics Port
    volumes:
      - nats_data:/data
    restart: always

volumes:
  nats_data:

Inisialisasi Stream dengan NATS CLI

Gunakan NATS CLI untuk membuat Stream bernama ORDERS. Stream ini menangkap semua event dengan subjek berawalan orders.>:

nats stream add ORDERS \
  --subjects "orders.>" \
  --storage file \
  --retention limits \
  --max-deliver 3 \
  --discard old \
  --dupe-window 2m

Analisis Parameter Konfigurasi Stream

  • --subjects "orders.>": Menggunakan wildcard > untuk menangkap semua hierarki subjek di bawah orders (contoh: orders.created, orders.paid, orders.cancelled).
  • --storage file: Pesan ditulis langsung ke disk untuk memastikan persistence mutlak saat server restart. Option memory dapat digunakan jika kecepatan in-memory menjadi prioritas utama.
  • --retention limits: Retention policy berbasis batas maksimum ukuran storage, umur pesan, atau jumlah pesan.
  • --max-deliver 3: Menginstruksikan broker untuk mencoba mengirim ulang pesan maksimal 3 kali jika terjadi failure atau timeout ACK.
  • --dupe-window 2m: Jendela waktu (2 menit) bagi broker untuk menyimpan fingerprint Nats-Msg-Id guna menolak duplikasi pesan.

3. Implementasi Microservice Producer & Worker di Golang

1. Inisialisasi Driver & Struct Schema

Pasang driver resmi NATS Go pada proyek Anda:

go get github.com/nats-io/nats.go

2. Implementasi Event Publisher (Producer)

Publisher bertugas mempublikasikan event bisnis ke NATS JetStream. Kode berikut mengimplementasikan deduplikasi menggunakan header nats.MsgId:

package main

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

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

type OrderCreatedEvent struct {
	OrderID     string    `json:"order_id"`
	CustomerID  string    `json:"customer_id"`
	Amount      float64   `json:"amount"`
	PaymentType string    `json:"payment_type"`
	CreatedAt   time.Time `json:"created_at"`
}

func main() {
	// Connect ke NATS Server
	nc, err := nats.Connect("nats://localhost:4222", nats.Timeout(5*time.Second))
	if err != nil {
		log.Fatalf("Gagal terhubung ke NATS: %v", err)
	}
	defer nc.Close()

	// Inisialisasi Context JetStream
	js, err := nc.JetStream()
	if err != nil {
		log.Fatalf("Gagal inisialisasi JetStream: %v", err)
	}

	event := OrderCreatedEvent{
		OrderID:     "ORD-89201",
		CustomerID:  "CUST-402",
		Amount:      750000.00,
		PaymentType: "CREDIT_CARD",
		CreatedAt:   time.Now(),
	}

	payload, err := json.Marshal(event)
	if err != nil {
		log.Fatalf("Gagal serialisasi JSON: %v", err)
	}

	// Publish event dengan MsgId unik untuk deduplikasi
	subject := "orders.created"
	pubAck, err := js.Publish(subject, payload, nats.MsgId(event.OrderID))
	if err != nil {
		log.Fatalf("Gagal publish pesan: %v", err)
	}

	fmt.Printf("Berhasil publish ke Stream: %s | Sequence: %d | Duplicate: %t\n",
		pubAck.Stream, pubAck.Sequence, pubAck.Duplicate)
}

3. Implementasi Durable Pull Consumer (Worker)

Pull Consumer merupakan pola terandal untuk microservice worker. Worker mengambil batch pesan sesuai kapasitas pemrosesan (backpressure control):

package main

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

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

type OrderCreatedEvent struct {
	OrderID     string    `json:"order_id"`
	CustomerID  string    `json:"customer_id"`
	Amount      float64   `json:"amount"`
	PaymentType string    `json:"payment_type"`
	CreatedAt   time.Time `json:"created_at"`
}

func main() {
	nc, err := nats.Connect("nats://localhost:4222")
	if err != nil {
		log.Fatalf("Connect fail: %v", err)
	}
	defer nc.Close()

	js, err := nc.JetStream()
	if err != nil {
		log.Fatalf("JetStream init fail: %v", err)
	}

	// Buat atau gunakan Durable Pull Subscription
	sub, err := js.PullSubscribe(
		"orders.created",
		"payment-service-worker",
		nats.Durable("payment-service-worker"),
		nats.ManualAck(),
	)
	if err != nil {
		log.Fatalf("Gagal buat subscription: %v", err)
	}

	log.Println("Worker Payment Service berjalan, menunggu event...")

	ctx := context.Background()
	for {
		// Fetch maksimal 10 pesan dalam window batch 5 detik
		msgs, err := sub.Fetch(10, nats.Context(ctx))
		if err != nil {
			// Timeout normal jika tidak ada pesan masuk
			continue
		}

		for _, msg := range msgs {
			processMessage(msg)
		}
	}
}

func processMessage(msg *nats.Msg) {
	var event OrderCreatedEvent
	if err := json.Unmarshal(msg.Data, &event); err != nil {
		log.Printf("Format payload tidak valid: %v", err)
		// Terminasi pesan buruk agar tidak masuk antrean retry
		_ = msg.Term()
		return
	}

	fmt.Printf("[Processing Order] ID: %s | Amount: %.2f\n", event.OrderID, event.Amount)

	// Simulasi proses transaksi bisnis
	err := executePayment(event)
	if err != nil {
		log.Printf("Proses pembayaran gagal untuk order %s: %v", event.OrderID, err)
		// Kembalikan ke stream untuk di-retry sesuai max_deliver policy
		_ = msg.Nak()
		return
	}

	// Konfirmasi sukses
	if err := msg.Ack(); err != nil {
		log.Printf("Gagal kirim ACK: %v", err)
	}
}

func executePayment(e OrderCreatedEvent) error {
	// Logic proses pembayaran
	return nil
}

4. Semantik Message Acknowledgement

JetStream menyediakan empat kontrol acknowledgement eksplisit:

  • msg.Ack(): Pemrosesan sukses. Broker menandai pesan sebagai dikonsumsi dan tidak akan mengirimkannya lagi ke consumer ini.
  • msg.Nak(): Pemrosesan gagal secara transient. Broker akan mengabaikan ack wait timeout dan langsung mengantrekan pesan untuk dikirim ulang.
  • msg.InProgress(): Digunakan untuk tugas berdurasi panjang (long-running job). Memberitahu broker bahwa consumer masih memproses pesan agar broker tidak menganggap pemrosesan timeout.
  • msg.Term(): Pemrosesan gagal secara permanen (malformed JSON, data corrupt). Broker menghentikan seluruh percobaan retry untuk pesan tersebut.

4. Best Practices Readiness Production

Mekanisme Deduplikasi Pesan

Sistem terdistribusi rentan terhadap pengiriman pesan ganda akibat network partition atau retry dari publisher. JetStream mengatasi masalah ini melalui fitur deduplikasi berbasis Nats-Msg-Id. Ketika publisher mengirim pesan dengan MsgId yang sudah ada di dalam dupe-window, broker menolak penulisan ulang ke stream log dan langsung mengembalikan ACK dengan status pubAck.Duplicate = true.

High Availability & Clustering (Raft Consensus)

Untuk lingkungan production, operasikan minimal 3 node NATS dalam satu cluster. JetStream menggunakan protokol konsensus Raft untuk mereplikasi log stream antar-node.

nats stream add ORDERS \
  --subjects "orders.>" \
  --storage file \
  --replicas 3

Dengan --replicas 3, data dipublikasikan ke node Leader, kemudian direplikasi secara konsisten ke dua node Follower sebelum ACK dikirim kembali ke publisher.

Monitoring, Metrics & Observability

NATS Server menyediakan endpoint metrik internal dalam format JSON dan HTTP Prometheus Exporter bawaan di port 8222:

  • http://localhost:8222/varz: Metrik umum engine NATS (CPU, RAM, Connections, In/Out bytes).
  • http://localhost:8222/jsz: Metrik spesifik JetStream (Stream status, storage usage, consumer lag, sequence numbers).

Perintah CLI untuk pemantauan realtime:

nats stream info ORDERS
nats consumer info ORDERS payment-service-worker
nats consumer report ORDERS

5. Checklist Readiness Production & Kesimpulan

Checklist Evaluasi Sistem

  • Cluster Setup: Minimal 3 node NATS dengan quorum Raft aktif.
  • Enkripsi Jaringan: TLS/mTLS wajib diaktifkan untuk lalu lintas client-to-server dan inter-node cluster.
  • Otentikasi & Otorisasi: Gunakan otentikasi berbasis NKEY atau JWT (NATS Account) untuk mengisolasi subjek antar-team.
  • Resource Limits: Tentukan max_bytes dan max_msgs per stream untuk mencegah kehabisan disk space storage.
  • Consumer Design: Selalu gunakan Pull Consumer dengan pembatasan Fetch() batch size sesuai resource pod worker Go.
  • Dead Letter Handling: Kombinasikan parameter max-deliver dengan logika msg.Term() pada worker untuk mengisolasi poison pill messages.

Kesimpulan

Integrasi NATS JetStream dan Golang menghasilkan fondasi mikroservis terdistribusi yang sangat cepat, efisien dari segi konsumsi memori, dan mudah dioperasikan. Berbeda dari Kafka atau RabbitMQ yang membutuhkan infrastruktur runtime rumit, NATS JetStream menyederhanakan arsitektur event-driven tanpa mengorbankan persistence data, deduplikasi, maupun keandalan pengiriman pesan.

๐Ÿ“– Artikel Terkait