Menjalankan task asynchronous pada aplikasi Go Fiber sering kali dimulai dengan goroutine biasa: go doTask(). Pola ini memicu masalah memory leak dan crash ketika beban request melonjak. Pilihan arsitektur kemudian mengerucut pada dua pendekatan: membangun in-process worker pool berbasis channel atau mengintegrasikan external message broker seperti Redis Streams, RabbitMQ, atau NATS.
Arsitektur Pemrosesan Latar Belakang: Internal Pool vs External Broker
In-process worker pool memanfaatkan Go runtime scheduler (GMP model). Aplikasi Go Fiber bertindak sekaligus sebagai HTTP server dan task consumer. Task dialirkan via Go channel buffered dan dieksekusi oleh pool goroutine yang jumlahnya dibatasi sejak inisialisasi.
Sebaliknya, arsitektur external broker memisahkan tanggung jawab penerimaan request (HTTP producer) dan eksekusi task (worker consumer). HTTP server hanya menerbitkan pesan ke antrean broker, kemudian instance worker terpisah mengonsumsi pesan tersebut secara asynchronous.
Analisis Trade-Off Teknis
1. CPU & Memory Contention terhadap HTTP Handling
Fiber berjalan di atas fasthttp yang dirancang untuk alokasi memori rendah dan event loop cepat. Jika worker internal mengeksekusi task intensif CPU (seperti hashing, manipulasi gambar, enkripsi payload besar) atau task I/O lambat, Go scheduler harus membagi waktu OS thread (P) antara HTTP handler dan worker goroutine.
Dampaknya langsung terasa pada latensi p99 HTTP request: request Fiber tertahan menunggu jatah scheduler. Pada broker eksternal, worker berjalan di pod atau container terpisah, sehingga lonjakan beban background task tidak memengaruhi performa HTTP handler.
2. Isolasi Kegagalan (Failure Isolation)
Worker internal berada dalam satu ruang memori dengan HTTP server. Satu task yang memicu panic tidak tertangani atau konsumsi memori tak terkontrol yang memicu Linux OOM (Out Of Memory) Killer akan mematikan seluruh proses Go Fiber. Akibatnya, traffic HTTP aktif langsung putus.
Broker eksternal menyediakan isolasi proses penuh. Jika worker mengalami crash atau OOM, pod worker tersebut di-restart oleh orchestrator (misal: Kubernetes) tanpa mengganggu koneksi pengguna di sisi API gateway atau web server.
3. Risiko Data Loss saat Graceful Shutdown
Saat aplikasi menerima sinyal terminasi (SIGTERM), task dalam in-memory channel berisiko hilang jika proses dihentikan paksa sebelum buffer kosong. Dibutuhkan koordinasi shutdown yang ketat antara HTTP server dan channel consumer.
Message broker eksternal memitigasi risiko ini lewat mekanisme acknowledgments (ACK/NACK) dan persistence ke disk. Pesan yang belum selesai diproses ketika worker mati akan di-requeue otomatis dan diambil oleh instance worker lain.
Perbandingan Biaya Cloud & Maintainability Codebase
Overhead RAM Pod vs Managed Broker
Secara finansial, in-process worker terlihat hemat karena tidak ada biaya managed service tambahan ($0 infra cost). Namun, pod aplikasi harus dialokasikan CPU dan RAM limit lebih tinggi (misalnya 1 vCPU / 1GB RAM per pod) untuk menampung lonjakan task mendadak. Autoscaling pod Fiber menjadi tidak efisien karena skala HTTP terikat dengan skala task processing.
Menggunakan managed broker (seperti AWS SQS, Upstash Redis, atau CloudAMQP) memunculkan biaya bulanan baseline (mulai $15 hingga $100+ per cluster). Namun, pod API Fiber dapat diperkecil ukurannya (misalnya 100m CPU / 128MB RAM) dan autoscaling worker diatur secara independen berdasarkan queue depth metric, bukan CPU load pod API.
Single Binary vs Distributed Service
In-process worker menjaga filosofi Go: single deployment artifact. Tidak perlu serialisasi kompleks, koneksi jaringan tambahan, service discovery, maupun maintenance kontrak skema antrean. Debugging cukup dengan tooling standar (pprof dan race detector).
Broker eksternal menuntut penanganan edge case sistem terdistribusi: penanganan dead-letter queue (DLQ), idempotency token, serialization overhead (JSON/Protobuf), serta monitoring broker lag.
Implementasi Minimalis: Worker Pool di Go Fiber
Contoh berikut mengimplementasikan in-process worker pool dengan graceful shutdown berbasis context.Context, penolakan beban saat queue penuh, dan tracking metrik runtime via sync/atomic.
package main
import (
"context"
"errors"
"log"
"os"
"os/signal"
"sync"
"sync/atomic"
"syscall"
"time"
"github.com/gofiber/fiber/v2"
)
type Task struct {
ID string
Payload string
}
type WorkerPool struct {
taskQueue chan Task
wg sync.WaitGroup
ctx context.Context
cancel context.CancelFunc
activeJobs int64
successJobs int64
}
func NewWorkerPool(workers int, queueSize int) *WorkerPool {
ctx, cancel := context.WithCancel(context.Background())
wp := &WorkerPool{
taskQueue: make(chan Task, queueSize),
ctx: ctx,
cancel: cancel,
}
for i := 0; i < workers; i++ {
wp.wg.Add(1)
go wp.worker()
}
return wp
}
func (wp *WorkerPool) worker() {
defer wp.wg.Done()
for {
select {
case <-wp.ctx.Done():
return
case task, ok := <-wp.taskQueue:
if !ok {
return
}
atomic.AddInt64(&wp.activeJobs, 1)
// ponytail: simulated payload processing; ganti dengan logic aktual
time.Sleep(100 * time.Millisecond)
atomic.AddInt64(&wp.activeJobs, -1)
atomic.AddInt64(&wp.successJobs, 1)
log.Printf("[Task Completed] ID: %s", task.ID)
}
}
}
func (wp *WorkerPool) Enqueue(task Task) error {
select {
case wp.taskQueue <- task:
return nil
default:
return errors.New("worker pool queue full")
}
}
func (wp *WorkerPool) Stop() {
wp.cancel()
close(wp.taskQueue)
wp.wg.Wait()
}
func main() {
pool := NewWorkerPool(4, 100)
app := fiber.New()
app.Post("/enqueue", func(c *fiber.Ctx) error {
taskID := c.Query("id", "default-id")
task := Task{ID: taskID, Payload: string(c.Body())}
if err := pool.Enqueue(task); err != nil {
return c.Status(fiber.StatusServiceUnavailable).JSON(fiber.Map{
"error": "Task queue full, try again later",
})
}
return c.Status(fiber.StatusAccepted).JSON(fiber.Map{
"status": "enqueued",
"task_id": taskID,
})
})
app.Get("/metrics", func(c *fiber.Ctx) error {
return c.JSON(fiber.Map{
"active_jobs": atomic.LoadInt64(&pool.activeJobs),
"completed_jobs": atomic.LoadInt64(&pool.successJobs),
"queue_depth": len(pool.taskQueue),
})
})
// Graceful shutdown handling
sigChan := make(chan os.Signal, 1)
signal.Notify(sigChan, os.Interrupt, syscall.SIGTERM)
go func() {
if err := app.Listen(":3000"); err != nil {
log.Printf("Server error: %v", err)
}
}()
<-sigChan
log.Println("Initiating graceful shutdown...")
_ = app.Shutdown()
pool.Stop()
log.Println("Shutdown complete.")
}
Dilewati: Retry policy otomatis, exponential backoff, persistent crash recovery. Tambahkan layer ini saat reliabilitas pesan mutlak diperlukan atau ganti dengan broker eksternal jika task wajib survive proses restart.
Matriks Keputusan: Kapan Harus Migrasi ke Broker?
Evaluasi opsi arsitektur menggunakan acuan kuantitatif berikut:
| Kriteria Evaluasi | In-Process Worker Pool | External Message Broker |
|---|---|---|
| Throughput Task | < 500 tasks/detik per pod | > 10.000 tasks/detik (skala cluster) |
| Toleransi Data Loss | Ada (non-kritis, e.g. cache warm, log sync) | Nol / Strict (e.g. order processing, payout) |
| Profil Beban Task | Ringan (I/O bound, HTTP webhook dispatch) | Berat (CPU bound, PDF rendering, video convert) |
| Latensi Toleransi HTTP | Fleksibel (p99 dapat terdistorsi saat peak) | Strict SLA (HTTP handler steril dari beban CPU task) |
| Ukuran Tim & Operasional | 1-3 engineer, prioritas kecepatan deploy | Tim menengah-besar, ada platform/DevOps support |
| Biaya Infrastruktur | Rendah pada skala awal, boros saat pod scale-up | Fixed cost managed broker, hemat pada skala besar |
Mulai dengan worker internal berukuran tetap jika task bersifat toleran terhadap kehilangan data dan tim masih mengutamakan kecepatan rilis. Segera beralih ke broker eksternal begitu p99 HTTP latency mulai terdegradasi akibat persaingan CPU atau saat integritas data task bernilai kritis bagi operasional bisnis.
Komentar
0 komentar
Masuk ke akun kamu untuk ikut berkomentar.
Belum ada komentar
Jadilah yang pertama ikut berdiskusi!