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:

  1. Worker A mengklaim lock untuk Task #101 dengan TTL 10 detik.
  2. 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.
  3. Pada detik ke-10, TTL lock di Redis habis. Lock dihapus secara otomatis oleh server.
  4. Queue orchestrator mendeteksi task belum selesai dan menugaskan kembali Task #101 ke Worker B.
  5. Worker B mengklaim lock baru untuk Task #101 dan mulai mengeksekusi komputasi.
  6. Pada detik ke-12, GC pause di Worker A selesai. Worker A tidak menyadari bahwa ia telah kehilangan kepemilikan lock.
  7. 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 computeBatch menerima sinyal workCtx.Done() dan menghentikan iterasi data tanpa mengeksekusi write.

Skenario: Latent Zombie Write

Jika worker mengalami OS-level freeze tepat sebelum memanggil w.storage.Commit:

  1. Lease habis di Redis. Worker B masuk, memperoleh fencingToken = 2, menyelesaikan kalkulasi, lalu menjalankan commit ke DB. Nilai last_fencing_token di DB menjadi 2.
  2. Worker A bangun kembali dari status freeze dengan fencingToken = 1.
  3. Worker A mengeksekusi SQL update dengan klausul WHERE last_fencing_token < 1.
  4. Database mengevaluasi ekspresi: kondisi 2 < 1 bernilai false. Row tidak terupdate.
  5. 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 / 3 dengan minimum TTL 5-10 detik.
  • Storage Prerequisite: Pola ini membutuhkan storage engine yang mendukung conditional writes atomic (misalnya SQL UPDATE ... WHERE, MongoDB findAndModify, 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.