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, blok defer, blok finally, 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 10000

Alih-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

  1. Buka transaksi database.
  2. 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;
  3. Cek rows affected. Jika bernilai 0, periksa status job saat ini:
    • Jika status COMPLETED: Lewati pemrosesan dan akui (ack) pesan antrean.
    • Jika status IN_PROGRESS tetapi updated_at sudah melampaui toleransi ambang batas (misal: > 10 menit), ambil alih proses (takeover) atau tandai sebagai stale.
  4. Jalankan mutasi data bisnis di dalam transaksi database yang sama.
  5. Perbarui status di processed_jobs menjadi COMPLETED lalu 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.