Akar Masalah: Batasan Shared Lock Konvensional
Pada distributed worker queue, race condition kerap terjadi saat worker mengeksekusi job kritis multi-tahap (seperti mutasi ledger keuangan atau provisioning infrastruktur). Pola umum menggunakan distributed lock berbasis Redis (seperti Redlock atau single-key TTL) rentan terhadap kegagalan konsistensi ketika terjadi:
- Stop-the-World GC Pause / Process Pause: Worker mengklaim lock dengan TTL 10 detik. Pada detik ke-9, worker mengalami garbage collection pause atau disk I/O stall selama 15 detik. TTL kedaluwarsa, queue server menganggap worker mati dan menjadwalkan ulang job ke Worker B. Worker A terbangun dan mengeksekusi commit ke database secara bersamaan dengan Worker B.
- Network Partition (Split-Brain): Worker terisolasi dari cluster queue tetapi masih memiliki akses ke database target. Worker baru ditugaskan untuk job yang sama oleh orchestrator.
- Asymmetric Network Delays: Heartbeat perpanjangan lease terlambat tiba di coordinator, memicu pelepasan sewa sepihak.
Mekanisme TTL lock standar tidak menjamin eksklusivitas mutual jika client dapat menahan eksekusi melewati durasi TTL tanpa disadari oleh downstream storage.
Prinsip Dasar: Quorum Lease dan Fencing Token
Pola Quorum Lease memecahkan masalah ini dengan memisahkan otoritas eksekusi menjadi dua komponen: validasi kuorum terdistribusi untuk kepemilikan lease, dan fencing token untuk verifikasi mutasi di sisi database target.
1. Verifikasi Kuorum (Raft/Paxos)
Lease tidak boleh diterbitkan oleh single node tanpa verifikasi mayoritas ($N/2 + 1$). State kepemilikan lease dikelola oleh sistem konsensus seperti etcd atau Consul. Worker hanya diizinkan menjalankan job jika lease tersebut masih divalidasi aktif oleh kuorum independen.
2. Monotonic Fencing Token (Epoch Increment)
Sesuai dengan analisis Martin Kleppmann mengenai distributed locking, lock server wajib menghasilkan token integer yang naik secara monotonik (epoch counter) setiap kali lease baru diberikan. Storage target (misalnya PostgreSQL atau MySQL) harus memvalidasi token ini dalam transaksi sebelum commit:
-- Skema tabel verifikasi epoch di database target
UPDATE critical_jobs
SET status = 'PROCESSING',
last_fencing_token = :current_token,
updated_at = NOW()
WHERE id = :job_id
AND last_fencing_token < :current_token;Jika Worker A terbangun dari GC pause dengan token 101, sedangkan Worker B telah mengambil alih dengan token 102, maka query update Worker A akan gagal karena kondisi last_fencing_token < 101 bernilai false.
Arsitektur Eksekusi Multi-Tahap
Job multi-tahap harus dipecah menjadi state machine deterministik di mana setiap transisi state diverifikasi terhadap kuorum lease:
- Acquire: Worker mengajukan lease ke etcd dengan TTL pendek (misal 5 detik) dan menerima fencing token generasi terbaru.
- Keep-Alive Background Loop: Goroutine/thread terpisah memperpanjang lease secara periodik (misal per 1.5 detik). Jika keep-alive gagal mencapai kuorum, sinyal pembatalan (context cancellation) langsung dikirim ke thread pemrosesan utama.
- Pre-Commit Check: Sebelum mutasi eksternal dijalankan, worker memverifikasi status lease lokal dan menyertakan token ke payload database.
- State Mutation: Database melakukan CAS (Compare-And-Swap) atomic pada kolom fencing token.
Implementasi Teknis: Worker Lease Berbasis etcd
Contoh implementasi teruji menggunakan Go dan primitif concurrency etcd v3:
package main
import (
"context"
"database/sql"
"errors"
"fmt"
"log"
"time"
clientv3 "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/concurrency"
)
type JobExecutor struct {
etcdCli *clientv3.Client
db *sql.DB
}
func (e *JobExecutor) ExecuteCriticalJob(ctx context.Context, jobID string) error {
// 1. Buat session dengan TTL 5 detik berbasis kuorum etcd
sess, err := concurrency.NewSession(e.etcdCli, concurrency.WithTTL(5))
if err != nil {
return fmt.Errorf("gagal membuat session kuorum: %w", err)
}
defer sess.Close()
// 2. Akuisisi mutex terdistribusi
lockKey := fmt.Sprintf("/leases/jobs/%s", jobID)
mutex := concurrency.NewMutex(sess, lockKey)
if err := mutex.Lock(ctx); err != nil {
return fmt.Errorf("gagal mengakuisisi lease: %w", err)
}
defer mutex.Unlock(context.Background())
// Header revision etcd berfungsi sebagai monotonic fencing token
fencingToken := mutex.Header().Revision
// 3. Jalankan tahap pekerjaan dengan validasi token di database
return e.processStageWithFencing(ctx, jobID, fencingToken)
}
func (e *JobExecutor) processStageWithFencing(ctx context.Context, jobID string, token int64) error {
tx, err := e.db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelReadCommitted})
if err != nil {
return err
}
defer tx.Rollback()
// Atomic CAS validation
res, err := tx.ExecContext(ctx, `
UPDATE job_states
SET current_epoch = $1, step = step + 1, updated_at = NOW()
WHERE id = $2 AND current_epoch < $1`,
token, jobID,
)
if err != nil {
return fmt.Errorf("error executing query: %w", err)
}
rowsAffected, err := res.RowsAffected()
if err != nil {
return err
}
if rowsAffected == 0 {
return errors.New("fencing token ditolak: lease telah kedaluwarsa atau diambil alih")
}
// Jalankan side effect di sini jika CAS berhasil
return tx.Commit()
}Mitigasi Node Flapping dan Unstable Network
Partisi jaringan intermiten dapat memicu node flapping, di mana worker berulang kali kehilangan dan merebut kembali lease. Mitigasi wajib diterapkan pada level scheduler:
- Exponential Backoff dengan Jitter: Saat lease terlepas di luar kendali normal, worker dilarang langsung melakukan re-acquire secara agresif. Terapkan full jitter:
sleep = min(cap, base * 2 ** attempt) * rand(0.5, 1.5). - Lease Renewal Grace Margin: Tetapkan batas renewal di level $1/3$ durasi TTL. Jika TTL 6 detik dan renewal tidak mendapat ack dalam 2 detik, tandai worker dalam kondisi degraded dan hentikan pengambilan job baru.
- Poison Pill Abort: Gunakan watchdog timer internal di runtime worker. Jika perpanjangan lease gagal mencapai kuorum sebelum 70% durasi TTL habis, worker mematikan prosesnya sendiri (fast-fail/SIGKILL) sebelum memicu race condition.
Strategi Safe Fallback Saat Kuorum Hilang
Bila koneksi ke cluster kuorum terputus di tengah pemrosesan:
- Seketika Batalkan Context: Seluruh I/O keluar (HTTP calls, DB queries) harus di-bind ke
context.Contextyang otomatis membatalkan operasi saat channelsess.Done()menerima sinyal. - Idempotency Verification: Setiap job wajib mengimplementasikan kunci idempotensi deterministik (misalnya hash input payload). Ketika worker baru mengambil alih job yang ditinggalkan, tahap yang sudah di-commit dengan token lama tidak akan dieksekusi ulang.
- Dead Letter Queue (DLQ) Routing: Jika sebuah job gagal diselesaikan setelah $N$ kali pergantian epoch (indikasi payload merusak worker/poison message), pindahkan task ke DLQ untuk audit manual.
Observabilitas dan Metrik Utama
Pantau metrik berikut untuk mendeteksi degradasi konsensus sebelum terjadi kegagalan split-brain:
| Metrik | Tipe | Batas Kritis / Anomali |
|---|---|---|
lease_acquisition_latency_ms | Histogram | P99 > 50% durasi heartbeat lease. Mengindikasikan disk lag pada etcd cluster. |
lease_drift_seconds | Gauge | Waktu sejak keep-alive terakhir yang berhasil. Nilai > 50% TTL memicu alert peringatan. |
quorum_rejection_total | Counter | Lonjakan nilai mengindikasikan worker stale yang mencoba menulis ke database (fencing aktif menolak commit). |
worker_evictions_total | Counter | Worker dipaksa terminate oleh internal watchdog karena kehilangan lease. |
Catatan Implementasi: Jangan pernah mengandalkan clock synchronization sistem (seperti NTP) untuk validasi lease di lingkungan terdistribusi. Clock skew dapat memicu drift waktu antar-node. Keamanan sistem harus murni bergantung pada monotonic sequence ordering via fencing token dan validasi kuorum Raft.
Komentar
0 komentar
Masuk ke akun kamu untuk ikut berkomentar.
Belum ada komentar
Jadilah yang pertama ikut berdiskusi!