Pada sistem komputasi batch terdistribusi, kegagalan worker sering diasumsikan sebagai crash-stop: proses mati total dan berhenti mengirim data. Asumsi ini keliru pada skenario beban komputasi tinggi. Beban CPU 100%, I/O thrashing, atau Garbage Collection (GC) stop-the-world pause sering menyebabkan worker berhenti merespons sementara waktu tanpa benar-benar mati.
Ketika distributed lock berbasis Time-To-Live (TTL) statis kadaluarsa saat proses sedang terhenti, antrean (message broker) menganggap worker mati dan menduplikasi task ke node lain. Saat worker pertama terbangun kembali, ia menjadi zombie worker yang tetap melanjutkan eksekusi dan menulis hasil ke storage. Kondisi ini menciptakan split-brain state dan merusak konsistensi data. Solusi deterministik untuk masalah ini membutuhkan kombinasi lease renewal, kill switch lokal, dan monotonic fencing token.
Anatomi Masalah: Mengapa TTL Statis Gagal
Pola distributed lock konvensional umumnya mengandalkan perintah atomik primitif seperti SET resource_key node_id NX PX 30000 di Redis. Pola ini mengasumsikan waktu eksekusi selalu berada di bawah durasi TTL.
Skenario kegagalan terjadi dalam kronologi berikut:
- Worker A mengklaim lock untuk Task #101 dengan TTL 10 detik.
- Worker A mulai memproses data batch besar. Pada detik ke-4, sistem operasi mengalami memory pressure tinggi yang memicu major GC pause selama 8 detik.
- Pada detik ke-10, TTL lock di Redis habis. Lock dihapus secara otomatis oleh server.
- Queue orchestrator mendeteksi task belum selesai dan menugaskan kembali Task #101 ke Worker B.
- Worker B mengklaim lock baru untuk Task #101 dan mulai mengeksekusi komputasi.
- Pada detik ke-12, GC pause di Worker A selesai. Worker A tidak menyadari bahwa ia telah kehilangan kepemilikan lock.
- Kedua worker mengeksekusi task yang sama secara simultan dan mencoba memodifikasi state akhir di database.
Peringatan: Mutex murni di sisi distributed cache/lock server tidak mampu membatasi akses klien jika klien tersebut mengalami uncoordinated pause. Komputasi terdistribusi tidak bisa bergantung pada clock synchronization atau asumsi durasi eksekusi tetap.
Tiga Komponen Proteksi
Untuk memastikan eksekusi aman dan dapat diinspeksi (inspectable), arsitektur worker membutuhkan tiga layer koordinasi:
- Heartbeat Lease Renewal: Routine latar belakang yang secara periodik memperpanjang TTL lock sebelum batas waktu habis, selama komputasi masih berjalan normal.
- Autonomous Kill Switch: Context cancellation lokal yang langsung memutus eksekusi komputasi bila heartbeat gagal memperpanjang lock sebelum safety threshold tercapai.
- Monotonic Fencing Token: Integer penomoran urut yang bertambah secara monotonik setiap kali lock baru diterbitkan. Komponen shared storage memvalidasi token ini dan menolak penulisan jika token yang dibawa klien lebih rendah daripada token terakhir yang telah diproses.
Implementasi Worker Loop dan Safety Controls
Contoh berikut mengimplementasikan worker batch dalam Go menggunakan model lease renewal asinkron dan context cancellation.
package main
import (
"context"
"errors"
"fmt"
"sync/atomic"
"time"
)
type Storage interface {
Commit(ctx context.Context, taskID string, fencingToken int64, result string) error
}
type LockManager interface {
AcquireLease(ctx context.Context, taskID string, ttl time.Duration) (token int64, err error)
RenewLease(ctx context.Context, taskID string, token int64, ttl time.Duration) error
ReleaseLease(ctx context.Context, taskID string, token int64) error
}
type BatchWorker struct {
lockMgr LockManager
storage Storage
ttl time.Duration
}
func (w *BatchWorker) ProcessTask(parentCtx context.Context, taskID string) error {
// 1. Klaim lease awal dan peroleh monotonic fencing token
fencingToken, err := w.lockMgr.AcquireLease(parentCtx, taskID, w.ttl)
if err != nil {
return fmt.Errorf("gagal klaim lock: %w", err)
}
defer w.lockMgr.ReleaseLease(context.Background(), taskID, fencingToken)
// Context lokal sebagai kill switch
workCtx, killSwitch := context.WithCancel(parentCtx)
defer killSwitch()
heartbeatErr := make(chan error, 1)
// 2. Heartbeat Goroutine untuk renewal lease
go func() {
ticker := time.NewTicker(w.ttl / 3)
defer ticker.Stop()
for {
select {
case <-workCtx.Done():
return
case <-ticker.C:
renewCtx, cancel := context.WithTimeout(workCtx, w.ttl/3)
err := w.lockMgr.RenewLease(renewCtx, taskID, fencingToken, w.ttl)
cancel()
if err != nil {
// Heartbeat gagal: matikan proses segera untuk mencegah zombie state
heartbeatErr <- err
killSwitch()
return
}
}
}
}()
// 3. Eksekusi komputasi batch
result, err := w.computeBatch(workCtx, taskID)
if err != nil {
select {
case hErr := <-heartbeatErr:
return fmt.Errorf("komputasi dibatalkan akibat lease renewal gagal: %w", hErr)
default:
return fmt.Errorf("komputasi gagal: %w", err)
}
}
// 4. Commit ke storage dengan membawa fencing token
if err := w.storage.Commit(parentCtx, taskID, fencingToken, result); err != nil {
return fmt.Errorf("commit ditolak: %w", err)
}
return nil
}
func (w *BatchWorker) computeBatch(ctx context.Context, taskID string) (string, error) {
// Simulasi komputasi bertahap dengan checkpoint context
for step := 1; step <= 5; step++ {
select {
case <-ctx.Done():
return "", ctx.Err()
case <-time.After(500 * time.Millisecond):
// Melakukan chunk processing
}
}
return "COMPUTE_SUCCESS", nil
}
Validasi Monotonic Fencing Token di Shared Storage
Heartbeat lokal tidak menjamin proteksi 100% jika worker mengalami freeze tepat sebelum baris commit dijalankan. Oleh karena itu, shared storage (PostgreSQL, MySQL, atau storage layer lainnya) wajib memberlakukan invariant fencing token.
Contoh skema dan atomic update pada PostgreSQL:
CREATE TABLE task_state (
task_id VARCHAR(64) PRIMARY KEY,
last_fencing_token BIGINT NOT NULL DEFAULT 0,
status VARCHAR(32) NOT NULL,
payload TEXT,
updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
);
-- Eksekusi commit oleh worker
UPDATE task_state
SET
status = 'COMPLETED',
payload = 'RESULT_DATA',
last_fencing_token = :incoming_token,
updated_at = NOW()
WHERE
task_id = :task_id
AND last_fencing_token < :incoming_token;
Jika query di atas menghasilkan 0 rows affected, storage driver worker harus menghasilkan error fatal. Hal ini menandakan worker lain dengan fencing token yang lebih tinggi (lebih baru) telah mengambil alih task atau telah menyelesaikan commit.
Skenario Partisi Jaringan dan Verifikasi Race Condition
Skenario: Silent Network Partition
Ketika kabel jaringan worker terputus hanya ke Redis/Lock Manager tetapi tetap terhubung ke Database:
- Ticker heartbeat mencoba melakukan renewal, namun mengalami
context.DeadlineExceeded. - Goroutine memanggil
killSwitch(). - Fungsi
computeBatchmenerima sinyalworkCtx.Done()dan menghentikan iterasi data tanpa mengeksekusi write.
Skenario: Latent Zombie Write
Jika worker mengalami OS-level freeze tepat sebelum memanggil w.storage.Commit:
- Lease habis di Redis. Worker B masuk, memperoleh
fencingToken = 2, menyelesaikan kalkulasi, lalu menjalankan commit ke DB. Nilailast_fencing_tokendi DB menjadi2. - Worker A bangun kembali dari status freeze dengan
fencingToken = 1. - Worker A mengeksekusi SQL update dengan klausul
WHERE last_fencing_token < 1. - Database mengevaluasi ekspresi: kondisi
2 < 1bernilai false. Row tidak terupdate. - Data hasil kalkulasi Worker B tetap utuh, mencegah split-brain corruption.
Trade-off dan Batasan Arsitektur
Pola lease lock dengan fencing token memiliki pertimbangan teknis berikut:
- Overhead Jaringan: Heartbeat interval yang terlalu rapat (misalnya di bawah 500ms) membebani Redis/lock engine jika terdapat puluhan ribu worker bersamaan. Aturan praktis yang stabil adalah
TTL / 3dengan minimum TTL 5-10 detik. - Storage Prerequisite: Pola ini membutuhkan storage engine yang mendukung conditional writes atomic (misalnya SQL
UPDATE ... WHERE, MongoDBfindAndModify, atau S3 Object Conditional Writes). Sistem file standar tanpa atomic lock tidak cocok digunakan sebagai target commit langsung. - Idempotensi Komputasi: Task yang dibatalkan oleh kill switch harus aman untuk diulang kembali oleh worker baru dari checkpoint terakhir.
Komentar
0 komentar
Masuk ke akun kamu untuk ikut berkomentar.
Belum ada komentar
Jadilah yang pertama ikut berdiskusi!