AizuDemy

Tutorial Build AI Agent Event-Driven Pakai LlamaIndex Workflows Python

Tutorial Build AI Agent Event-Driven Pakai LlamaIndex Workflows Python
๐ŸŽง
Dengarkan Artikel Ini
Suara AI Otomatis โ€ข 7 mnt baca baca
โšก TL;DR

Poin Kunci Artikel Ini:

  • Tangkap exception pada step handler yang memanggil LLM.
  • Evaluasi counter retry pada event payload.
  • Limitasi Agent Synchronous dan Urgensi Event-Driven ArchitectureArsitektur AI Agent tradisional umumnya dibangun menggunakan eksekusi synchronous linear.
๐Ÿ“‹ Daftar Isi Materi Tutup โ–ด

1. Limitasi Agent Synchronous dan Urgensi Event-Driven Architecture

Arsitektur AI Agent tradisional umumnya dibangun menggunakan eksekusi synchronous linear. Pada pola ini, sebuah langkah (Step A) memanggil LLM API, menghentikan eksekusi thread (blocking I/O) saat menunggu respon HTTP, lalu melanjutkan ke Langkah B. Pendekatan linear ini menimbulkan kelemahan struktural serius saat sistem harus menangani workload skala besar di tingkat produksi.

Masalah utama agent synchronous melingkupi:

  • Thread Exhaustion: Setiap request mengunci thread OS atau event loop I/O. Ketika latency LLM meningkat dari 500ms menjadi 10 detik, thread pool pada application server (seperti Gunicorn atau Uvicorn) akan habis dengan cepat.
  • Ketidakmampuan Menangani Task Asentris: Agent modern butuh mengeksekusi multiple sub-task secara paralel (misal: merangkum dokumen, memanggil API eksternal, dan mengeksekusi pencarian vektor bersamaan). Pola synchronous memaksa sub-task berjalan berurutan, melipatgandakan total waktu eksekusi (latency akumulatif).
  • Resilience Buruk: Jika salah satu panggilan LLM di tengah alur mengalami timeout, seluruh execution state hilang kecuali diimplementasikan mekanisme persistence manual yang rumit.

Solusinya adalah mentransformasi arsitektur menjadi Event-Driven State Machine (EDSM). Dalam arsitektur event-driven, komponen tidak saling memanggil secara langsung. Sebaliknya, setiap komponen (step) mengonsumsi event, melakukan kalkulasi async, dan memancarkan (emit) event baru ke event bus. LlamaIndex Workflows menyediakan abstraksi primitif tingkat tinggi native Python untuk membangun arsitektur ini tanpa overhead boilerplate code yang berat.

2. Primitif Utama LlamaIndex Workflows

LlamaIndex Workflows mengabstraksi eksekusi agent menggunakan lima komponen fondasi:

  1. Workflow: Kelas dasar pengatur event loop internal, registrasi step, dan manajemen lifecycle alur kerja.
  2. Event: Data Transfer Object (DTO) yang membawa payload state antar-step. Setiap custom event merupakan subclass dari llama_index.core.workflow.Event.
  3. @step: Python decorator untuk mendaftarkan fungsi async sebagai handler event tertentu.
  4. StartEvent dan StopEvent: Event bawaan sistem penanda titik masuk (entrypoint) dan titik keluar (exitpoint) aliran data.
  5. Context: Objek manajemen state global yang aman diakses secara ko-kuren oleh berbagai step untuk menyimpan variabel bersama.

3. Persiapan Lingkungan Pengembangan

Siapkan virtual environment Python versi 3.10 atau lebih baru. Install package core LlamaIndex beserta integrasi LLM Ollama untuk eksekusi lokal:

python3 -m venv venv
source venv/bin/activate
pip install llama-index-core llama-index-llms-ollama ollama

Pastikan runtime Ollama terinstal pada sistem dan unduh model LLM yang akan digunakan:

ollama pull llama3.2

4. Implementasi Workflow Multi-Step dengan Pola Fan-Out / Fan-In

Berikut adalah implementasi sistem agent event-driven lengkap. Kode ini menunjukkan pola Fan-Out (meneruskan satu event ke multiple step paralel) dan Fan-In (menggabungkan hasil dari beberapa step sebelum menghasilkan output final):

import asyncio
from typing import Any
from llama_index.core.workflow import (
    Workflow,
    Event,
    StartEvent,
    StopEvent,
    step,
    Context
)
from llama_index.llms.ollama import Ollama

class QueryEvent(Event):
    query: str

class SentimentResultEvent(Event):
    sentiment: str

class SummaryResultEvent(Event):
    summary: str

class ParallelAgentWorkflow(Workflow):
    def __init__(self, llm: Ollama, **kwargs: Any):
        super().__init__(**kwargs)
        self.llm = llm

    @step
    async def setup_query(self, ctx: Context, ev: StartEvent) -> QueryEvent:
        user_input = ev.get("user_input")
        if not user_input:
            raise ValueError("User input tidak boleh kosong.")
        await ctx.set("raw_query", user_input)
        return QueryEvent(query=user_input)

    @step
    async def analyze_sentiment(self, ev: QueryEvent) -> SentimentResultEvent:
        prompt = f"Analisis sentimen teks berikut singkat saja (Positif/Negatif/Netral): {ev.query}"
        response = await self.llm.acomplete(prompt)
        return SentimentResultEvent(sentiment=str(response).strip())

    @step
    async def summarize_text(self, ev: QueryEvent) -> SummaryResultEvent:
        prompt = f"Buat rangkuman maksimal 2 kalimat dari teks berikut: {ev.query}"
        response = await self.llm.acomplete(prompt)
        return SummaryResultEvent(summary=str(response).strip())

    @step
    async def aggregate_results(
        self, ctx: Context, ev: SentimentResultEvent | SummaryResultEvent
    ) -> StopEvent | None:
        # Simpan event yang masuk ke context
        if isinstance(ev, SentimentResultEvent):
            await ctx.set("sentiment", ev.sentiment)
        elif isinstance(ev, SummaryResultEvent):
            await ctx.set("summary", ev.summary)

        # Cek apakah kedua data paralel sudah selesai diproses
        sentiment = await ctx.get("sentiment", default=None)
        summary = await ctx.get("summary", default=None)

        if sentiment is not None and summary is not None:
            final_output = {
                "original_query": await ctx.get("raw_query"),
                "sentiment": sentiment,
                "summary": summary
            }
            return StopEvent(result=final_output)
        
        return None

async def main():
    llm = Ollama(model="llama3.2", request_timeout=120.0)
    workflow = ParallelAgentWorkflow(llm=llm, timeout=60, verbose=True)
    
    input_text = (
        "Sistem cloud baru kami mengalami downtime selama 2 jam tadi malam karena kesalahan konfigurasi DNS. "
        "Namun, tim insinyur merespons dengan sangat cepat dan berhasil memulihkan seluruh layanan tanpa kehilangan data pelanggan."
    )
    
    print("Memulai eksekusi event-driven workflow...")
    result = await workflow.run(user_input=input_text)
    print("\nHasil Eksekusi Aggregator:")
    print(result)

if __name__ == "__main__":
    asyncio.run(main())

Detail Mekanisme Kerja Kode:

  • Saat workflow.run() dipanggil, StartEvent dipancarkan ke event loop.
  • Step setup_query menangkap StartEvent dan memancarkan QueryEvent.
  • Step analyze_sentiment dan summarize_text mendengarkan tipe event yang sama (QueryEvent). Kedua step ini dieksekusi secara terpisah dan paralel tanpa saling menunggu.
  • Step aggregate_results bertindak sebagai pengumpul (Fan-In). Step ini dipanggil setiap kali salah satu dari SentimentResultEvent atau SummaryResultEvent dipancarkan. Step mengecek state pada Context; jika kedua data belum lengkap, step mengembalikan None untuk tetap menunggu. Setelah seluruh data lengkap, step memancarkan StopEvent yang mengakhiri eksekusi workflow.

5. Resilience, Error Handling, Retry Pattern, dan Pengelolaan Memori

Integrasi LLM pada lingkungan produksi rentan terhadap kegagalan transient (seperti HTTP 503, rate limit 429, atau socket timeout). Menggunakan block try-except sederhana di dalam synchronous function tidak cukup. Workflow butuh Retry Pattern berbasis event.

Urutan Eksekusi Error Handling Penanganan Retry:

  1. Tangkap exception pada step handler yang memanggil LLM.
  2. Evaluasi counter retry pada event payload.
  3. Jika retry belum melebihi batas maksimum (misal 3 kali), pancarkan RetryEvent baru yang membawa informasi attempt counter dan delay waktu (Exponential Backoff).
  4. Step khusus handle_retry menerima RetryEvent, melakukan asyncio.sleep() sesuai delay, lalu memancarkan ulang event pemrosesan awal.
  5. Jika batas retry terlampaui, alihkan ke ErrorEvent atau kembalikan StopEvent berisi payload error terstruktur.

Implementasi retry pattern terstruktur:

import asyncio
from llama_index.core.workflow import Workflow, Event, StartEvent, StopEvent, step

class LLMCallEvent(Event):
    prompt: str
    retries: int = 0

class ResilientWorkflow(Workflow):
    @step
    async def process_call(self, ev: StartEvent | LLMCallEvent) -> StopEvent | LLMCallEvent:
        if isinstance(ev, StartEvent):
            prompt = ev.get("prompt")
            retries = 0
        else:
            prompt = ev.prompt
            retries = ev.retries

        try:
            # Simulasi eksekusi LLM yang dapat gagal
            res = await self.llm.acomplete(prompt)
            return StopEvent(result=str(res))
        except Exception as err:
            if retries < 3:
                backoff_delay = 2 ** retries
                print(f"[Error] {err}. Retry ke-{retries + 1} dalam {backoff_delay} detik...")
                await asyncio.sleep(backoff_delay)
                return LLMCallEvent(prompt=prompt, retries=retries + 1)
            
            return StopEvent(result=f"Execution Failed permanently: {str(err)}")

Manajemen Memori pada Workflow Berumur Panjang (Long-Running):

Event payload dan objek Context disimpan dalam RAM selama instance workflow berjalan. Jika workflow memproses file besar (seperti RAG berbasis PDF ratusan halaman), perhatikan pencegahan memory leak berikut:

  • Payload Scoping: Hindari memasukkan objek besar (misal: instance Document atau tensor embedding) ke dalam kelas Event. Sertakan hanya identifier (misal: ID dokumen atau URI storage). Letakkan data besar pada external store seperti Redis atau S3.
  • Purge Context Variable: Hapus variabel context yang tidak lagi dibutuhkan menggunakan await ctx.set("key", None) setelah step konsumsi selesai.
  • Garbage Collection Explicit: Pada daemon worker yang memproses ribuan event per menit, panggil gc.collect() secara periodik untuk membebaskan unreferenced circular references.

6. Human-in-the-Loop (HITL) dan Streaming Event

Workflow event-driven memungkinkan penghentian sementara (pause) eksekusi untuk menunggu persetujuan manusia (Human-in-the-Loop) atau streaming status secara real-time ke UI frontend.

Untuk mengimplementasikan streaming status internal workflow ke client (misalnya via Server-Sent Events/SSE pada HTTP API), gunakan event handler bawaan workflow.stream_events():

async def run_with_streaming():
    wf = ParallelAgentWorkflow(llm=llm)
    handler = wf.run(user_input="Input teks analisis...")
    
    # Stream event internal yang dipancarkan oleh step
    async for event in wf.stream_events():
        print(f"[Event Emitted]: {type(event).__name__}")
        
    final_result = await handler
    print("Hasil akhir:", final_result)

7. Deploy Ke Lingkungan Production

Saat men-deploy LlamaIndex Workflows ke lingkungan produksi, arsitektur harus memenuhi standar high availability dan observability.

Checklist Production Ready:

  • Asynchronous Entrypoint: Bungkus workflow di dalam HTTP framework non-blocking seperti FastAPI menggunakan ASGI server (Uvicorn / Hypercorn). Hindari WSGI server seperti Flask synchronous.
  • Distributed Task Execution: Jika workflow membutuhkan koordinasi multi-node, integrasikan event engine dengan Redis Pub/Sub, RabbitMQ, atau Temporal.io sebagai transport layer eksternal.
  • Strict Timeout Management: Selalu atur `timeout` eksplisit pada konstruktor `Workflow(timeout=30.0)` serta pada LLM client untuk mencegah worker hanging selamanya.
  • Observability dan Tracing: Hubungkan workflow dengan platform OpenTelemetry atau Arize Phoenix. Phoenix menyediakan tracing otomatis untuk melihat grafik latensi per-step, jejak eksekusi event, serta payload LLM secara transparan.

๐Ÿ“– Artikel Terkait