Akar Masalah: Variansi Latensi Ekstrem pada Pipeline AI

Pipeline AI code generation menghadapi variansi latensi yang lebar. Task klasifikasi prompt, auto-complete, atau unit test patch selesai dalam rentang 500 milidetik hingga 2 detik. Sebaliknya, task sintesis multi-file, refactoring arsitektur, atau code generation berbasis ReAct loop membutuhkan waktu antara 45 hingga 180 detik.

Ketika kedua karakteristik beban kerja ini dicampur ke dalam antrean FIFO (First-In, First-Out) tunggal, timbul anomali Head-of-Line (HoL) blocking. Beberapa task sintesis kode berdurasi panjang menempati seluruh worker concurrency slot yang tersedia. Akibatnya, task interaktif ringan tertahan di belakang antrean, menaikkan p99 latency secara dramatis dan memicu worker starvation pada layanan interaktif.

Arsitektur Segregasi Antrean dan Dynamic Allocation

Solusi pertama adalah memisahkan antrean berdasarkan bobot komputasi dan SLA waktu respon:

  • Antrean interactive: Untuk task berlatensi rendah (< 3 detik) seperti auto-complete, code explanation, dan lint patch.
  • Antrean synthesis: Untuk task berdurasi panjang (> 10 detik) seperti scaffold project, multi-file code generation, dan agent execution loop.

Alokasi worker tidak boleh dibagi secara statis 50:50. Pembagian kaku menyebabkan pemborosan resource saat salah satu antrean sepi. Gunakan weighted dynamic consumer pool:

+-----------------------+      +-----------------------+
|  Client (API Gateway) |      | Client (Batch / Plan) |
+-----------+-----------+      +-----------+-----------+
            |                              |
            v                              v
   [queue:interactive]             [queue:synthesis]
            |                              |
            +--------------+---------------+
                           |
            +--------------v---------------+
            |    Dynamic Worker Pool       |
            | - 70% slot reserved: interac |
            | - 30% slot reserved: synth   |
            | - Idle workers steal tasks   |
            +------------------------------+

Problem Visibility Timeout dan Solusi Heartbeat Lease Lock

Sistem antrean pesan seperti RabbitMQ, SQS, atau Redis Streams mengandalkan ack timeout atau visibility timeout. Jika worker tidak mengirim ack sebelum timeout habis, antrean menganggap worker mati dan mendistribusikan ulang task ke worker lain. Pada pipeline LLM yang memakan waktu 90 detik, menyetel timeout statis ke 120 detik membawa risiko besar:

  • Bila worker benar-benar crash di detik ke-5, task tertahan 115 detik sebelum di-retry.
  • Bila sintesis melampaui 120 detik karena LLM provider throttled, task di-dispatch ke worker kedua, menghasilkan duplicate code synthesis dan pemborosan biaya token API ganda.

Solusinya adalah Heartbeat Lease Lock. Set visibility timeout pendek (misal: 15 detik), lalu jalankan background thread/task pada worker yang secara berkala memperpanjang lease (TTL) di Redis selama proses LLM masih berjalan aktif.

Implementasi Heartbeat Lease Lock dan Idempotency

Berikut implementasi heartbeat lock berbasis Redis dan eksekusi task idempoten menggunakan Python:

import asyncio
import time
import uuid
import redis.asyncio as redis

class AIJobWorker:
    def __init__(self, redis_client: redis.Redis):
        self.redis = redis_client
        # ponytail: implementasi in-process heartbeat lock. Tingkatkan ke Redlock jika cluster multi-master.
        self._renew_script = self.redis.register_script("""
            if redis.call('get', KEYS[1]) == ARGV[1] then
                return redis.call('pexpire', KEYS[1], ARGV[2])
            else
                return 0
            end
        """)

    async def _heartbeat(self, lock_key: str, lock_val: str, interval_ms: int, ttl_ms: int, stop_event: asyncio.Event):
        while not stop_event.is_set():
            try:
                await asyncio.sleep(interval_ms / 1000.0)
                if stop_event.is_set():
                    break
                result = await self._renew_script(keys=[lock_key], args=[lock_val, ttl_ms])
                if result == 0:
                    # Lock terlepas atau diambil alih worker lain
                    stop_event.set()
                    break
            except Exception:
                # Log kegagalan network, coba kembali pada loop berikutnya jika belum timeout
                pass

    async def execute_task(self, idempotency_key: str, task_fn, *args, **kwargs):
        lock_key = f"lock:task:{idempotency_key}"
        lock_val = str(uuid.uuid4())
        ttl_ms = 15000  # 15 detik TTL
        renew_interval_ms = 5000  # Kirim detak jantung per 5 detik

        # 1. Akuisisi Lock secara eksklusif (NX)
        acquired = await self.redis.set(lock_key, lock_val, px=ttl_ms, nx=True)
        if not acquired:
            # Task sedang berjalan di worker lain atau duplikat
            return None

        stop_event = asyncio.Event()
        heartbeat_task = asyncio.create_task(
            self._heartbeat(lock_key, lock_val, renew_interval_ms, ttl_ms, stop_event)
        )

        try:
            # 2. Eksekusi task code generation berdurasi panjang
            result = await task_fn(*args, **kwargs)
            return result
        finally:
            # 3. Hentikan heartbeat dan rilis lock
            stop_event.set()
            await heartbeat_task
            release_script = """
                if redis.call('get', KEYS[1]) == ARGV[1] then
                    return redis.call('del', KEYS[1])
                else
                    return 0
                end
            """
            await self.redis.eval(release_script, 1, lock_key, lock_val)

Idempotency di Batas Sistem (Ingress Layer)

Selain worker-level lock, cegah antrean menerima duplikasi sejak entrypoint API. Buat hash deterministik dari parameter payload:

import hashlib
import json

def generate_idempotency_key(project_id: str, prompt: str, target_files: list) -> str:
    payload = {
        "project_id": project_id,
        "prompt": prompt.strip(),
        "target_files": sorted(target_files)
    }
    encoded = json.dumps(payload, sort_keys=True).encode("utf-8")
    return hashlib.sha256(encoded).hexdigest()

Sebelum payload dimasukkan ke queue broker, periksa status key tersebut di Redis:

  1. Jika key ada dan status PENDING atau PROCESSING, kembalikan HTTP 202 beserta Task ID yang sama.
  2. Jika status COMPLETED, langsung kembalikan output dari storage cache.
  3. Jika key belum ada, set status PENDING dengan TTL wajar (misal: 1 jam), kemudian LPUSH atau kirim event ke message broker.

Ringkasan Trade-Off dan Troubleshooting

Penerapan pola ini menyelesaikan masalah starvation, tetapi membawa sejumlah konsekuensi teknis:

  • Complexity Overhead: Mengelola background heartbeat task mewajibkan penanganan koneksi async yang bersih. Pastikan background loop ditutup via blok finally agar tidak memicu memory leak.
  • Deadlock Prevention: TTL pendek (15 detik) menjamin bahwa jika worker mati mendadak (OOMKilled atau node failure), antrean tidak terkunci permanen. Lease akan kadaluarsa secara alami dalam 15 detik, sehingga mekanisme retry dapat mengambil alih tugas tanpa intervensi manual.
  • Queue Starvation Metrik: Pantau metrik queue_depth dan task_latency_seconds terpisah untuk setiap segmen queue. Lonjakan depth pada queue interactive menjadi sinyal langsung untuk menaikkan horizontal autoscaler pada reserved interactive worker pool.