Anatomi Kegagalan: Mengapa Worker Menjadi Zombie dan Lock Menjadi Basi
Dalam sistem pemrosesan asinkron skala produksi, kegagalan worker adalah keniscayaan. Masalah fatal terjadi ketika worker berhenti memproses tugas namun statusnya di sistem antrean tetap menggantung (job zombie). Tiga penyebab utamanya meliputi:
- Linux Out-Of-Memory (OOM) Killer: Kernel menghentikan proses worker menggunakan sinyal
SIGKILL(sinyal 9). Proses mati seketika tanpa mengeksekusi penanganan sinyal, blokdefer, blokfinally, maupun pelepasan lock. - I/O Blocking Tanpa Deadline: Permintaan HTTP pihak ketiga atau query database mengalami TCP connection hang tanpa timeout ketat. Thread tetap hidup, tetapi tidak pernah menyelesaikan eksekusi (hung state).
- Network Partition Flapping: Koneksi worker ke broker terputus sesaat setelah pesan diambil (di-ack), meninggalkan state transaksi tidak jelas.
Jika sistem menggunakan distributed lock statis tanpa kedaluwarsa atau dengan durasi manual, resource yang terkunci tidak akan pernah dilepas. Akibatnya antrean downstream macet (deadlock) dan membutuhkan intervensi manual developer untuk menghapus kunci di penyimpanan data.
Heartbeat Lease Renewal dan TTL Auto-Eviction
Pola static distributed lock tidak memadai untuk sistem mandiri (self-healing). Mekanisme yang benar adalah menggunakan sewa berbasis waktu (lease) dengan pembaruan berkala (heartbeat) dan auto-eviction melalui TTL.
1. Redis Distributed Lock dengan Lease
Kunci diakuisisi dengan parameter NX (set if not exists) dan PX (milliseconds expiration) menggunakan identitas unik token per worker instance:
SET resource_key <unique_worker_token> NX PX 10000Alih-alih mengandalkan estimasi durasi job yang rentan keliru, tetapkan TTL pendek (misalnya 10 detik). Jalankan goroutine atau background thread heartbeat yang memperpanjang TTL (misalnya setiap 3 detik) menggunakan evaluasi script Lua selama proses berjalan normal.
Jika worker mengalami crash akibat OOM, thread heartbeat berhenti. Redis secara otomatis menghapus kunci setelah 10 detik kedaluwarsa (auto-eviction), memungkinkan worker lain mengambil alih tugas tanpa intervensi manusia.
2. PostgreSQL Advisory Lock vs Redis Lease
PostgreSQL menyediakan fungsi advisory lock seperti pg_try_advisory_xact_lock(bigint):
- PostgreSQL Advisory Lock: Mengikat siklus hidup lock pada koneksi TCP atau transaksi aktif. Jika worker crash, PostgreSQL segera mendeteksi koneksi putus dan melepaskan lock secara otomatis. Keterbatasannya: membebani connection pool database utama.
- Redis Lease: Skalabel untuk throughput ribuan transaksi per detik tanpa menguras pool transaksi database relasional, namun memerlukan validasi token saat perpanjangan dan pelepasan lock guna menghindari lock hijacking akibat race condition.
Implementasi Heartbeat Lock dan Context Timeout di Go
Contoh berikut mendemonstrasikan eksekusi worker di Go dengan context timeout, pelepasan kunci atomik via Lua script, dan perpanjangan lease periodik:
package main
import (
"context"
"crypto/rand"
"encoding/hex"
"errors"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
var (
renewScript = redis.NewScript(`
if redis.call("get", KEYS[1]) == ARGV[1] then
return redis.call("pexpire", KEYS[1], ARGV[2])
else
return 0
end
`)
releaseScript = redis.NewScript(`
if redis.call("get", KEYS[1]) == ARGV[1] then
return redis.call("del", KEYS[1])
else
return 0
end
`)
)
type AutoRenewLock struct {
client *redis.Client
key string
token string
ttl time.Duration
}
func AcquireLock(ctx context.Context, rdb *redis.Client, key string, ttl time.Duration) (*AutoRenewLock, error) {
b := make([]byte, 16)
if _, err := rand.Read(b); err != nil {
return nil, err
}
token := hex.EncodeToString(b)
ok, err := rdb.SetNX(ctx, key, token, ttl).Result()
if err != nil {
return nil, err
}
if !ok {
return nil, errors.New("gagal mengakuisisi lock: sedang digunakan worker lain")
}
return &AutoRenewLock{
client: rdb,
key: key,
token: token,
ttl: ttl,
}, nil
}
func (l *AutoRenewLock) StartHeartbeat(ctx context.Context, interval time.Duration) {
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
ttlMs := int(l.ttl.Milliseconds())
res, err := renewScript.Run(context.Background(), l.client, []string{l.key}, l.token, ttlMs).Result()
if err != nil || res.(int64) == 0 {
return
}
}
}
}()
}
func (l *AutoRenewLock) Release(ctx context.Context) error {
_, err := releaseScript.Run(ctx, l.client, []string{l.key}, l.token).Result()
return err
}
func ProcessJob(parentCtx context.Context, rdb *redis.Client, jobID string) error {
jobCtx, cancel := context.WithTimeout(parentCtx, 45*time.Second)
defer cancel()
lockKey := fmt.Sprintf("lock:job:%s", jobID)
ttl := 10 * time.Second
lock, err := AcquireLock(jobCtx, rdb, lockKey, ttl)
if err != nil {
return fmt.Errorf("lewati job %s: %w", jobID, err)
}
defer func() {
relCtx, relCancel := context.WithTimeout(context.Background(), 2*time.Second)
defer relCancel()
_ = lock.Release(relCtx)
}()
lock.StartHeartbeat(jobCtx, 3*time.Second)
return executeTask(jobCtx)
}
func executeTask(ctx context.Context) error {
select {
case <-time.After(2 * time.Second):
return nil
case <-ctx.Done():
return ctx.Err()
}
}
Retry Resilience: Exponential Backoff, Jitter, dan DLQ
Mengulang eksekusi (retry) job gagal secara agresif tanpa penundaan akan menyebabkan thundering herd problem yang melumpuhkan database atau downstream API yang sedang pulih. Gunakan formula Full Jitter:
Sleep = random(0, min(MaxBackoff, BaseBackoff * (2 ^ attempt)))import (
"math"
"math/rand"
"time"
)
func CalculateBackoff(attempt int, base time.Duration, maxBackoff time.Duration) time.Duration {
multiplier := math.Pow(2, float64(attempt))
temp := float64(base) * multiplier
if temp > float64(maxBackoff) {
temp = float64(maxBackoff)
}
return time.Duration(rand.Float64() * temp)
}
Routing ke Dead Letter Queue (DLQ)
Job diarahkan ke antrean khusus (DLQ) jika memenuhi salah satu kondisi berikut:
- Jumlah percobaan melampaui batas maksimal (misal:
attempt > 5). - Error diklasifikasikan sebagai unrecoverable atau poison pill (misalnya: data korup, validasi payload gagal, atau record tidak ditemukan). Error tipe ini tidak boleh diretry.
DLQ mengisolasi pesan bermasalah tanpa memblokir pemrosesan antrean utama, sehingga metrik alerting dapat dipicu secara spesifik.
Idempotent Consumer Melalui Deduplication Table
Saat terjadi auto-eviction atau kegagalan jaringan setelah proses bisnis selesai tetapi sebelum lock dirilis, sistem antrean berpotensi menjalankan ulang job yang sama (at-least-once delivery). Tanpa penanganan idempotensi, eksekusi ulang memicu mutasi ganda seperti double payment.
Gunakan deduplication table di database transaksional dengan memanfaatkan integritas UNIQUE CONSTRAINT:
CREATE TABLE processed_jobs (
idempotency_key VARCHAR(128) PRIMARY KEY,
job_name VARCHAR(64) NOT NULL,
status VARCHAR(32) NOT NULL,
response_payload JSONB,
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);
Alur Eksekusi Idempoten
- Buka transaksi database.
- Jalankan perintah insert idempotency key:
INSERT INTO processed_jobs (idempotency_key, job_name, status) VALUES ($1, $2, 'IN_PROGRESS') ON CONFLICT (idempotency_key) DO NOTHING; - Cek rows affected. Jika bernilai
0, periksa status job saat ini:- Jika status
COMPLETED: Lewati pemrosesan dan akui (ack) pesan antrean. - Jika status
IN_PROGRESStetapiupdated_atsudah melampaui toleransi ambang batas (misal: > 10 menit), ambil alih proses (takeover) atau tandai sebagai stale.
- Jika status
- Jalankan mutasi data bisnis di dalam transaksi database yang sama.
- Perbarui status di
processed_jobsmenjadiCOMPLETEDlalu commit transaksi.
Metrik Observabilitas Penting
Sistem self-healing memerlukan telemetri yang jelas agar tim dapat mengukur anomali tanpa harus masuk ke server secara manual:
- Queue Lag & Queue Depth: Selisih antara timestamp job dibuat dengan waktu job mulai diproses. Peningkatan lag mendadak menandakan kapasitas worker tidak sebanding dengan laju pesan masuk.
- Lock Age: Durasi aktif suatu lock. Jika lock mendekati ambang batas context timeout, identifikasi apakah dependensi downstream mengalami penurunan performa (degradasi).
- Heartbeat Failure Rate: Metrik penghitung kegagalan evaluasi Lua script heartbeat. Peningkatan nilai ini menunjukkan adanya degradasi koneksi antara worker dan cluster Redis.
- DLQ Ingestion Counter: Jumlah job yang masuk ke antrean DLQ. Pemicu alarm utama bagi tim teknis untuk menganalisis payload error yang membutuhkan patch kode.
Komentar
0 komentar
Masuk ke akun kamu untuk ikut berkomentar.
Belum ada komentar
Jadilah yang pertama ikut berdiskusi!