Tutorial Build AI Agent Event-Driven Pakai LlamaIndex Workflows Python
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.
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:
Workflow: Kelas dasar pengatur event loop internal, registrasi step, dan manajemen lifecycle alur kerja.Event: Data Transfer Object (DTO) yang membawa payload state antar-step. Setiap custom event merupakan subclass darillama_index.core.workflow.Event.@step: Python decorator untuk mendaftarkan fungsi async sebagai handler event tertentu.StartEventdanStopEvent: Event bawaan sistem penanda titik masuk (entrypoint) dan titik keluar (exitpoint) aliran data.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 ollamaPastikan runtime Ollama terinstal pada sistem dan unduh model LLM yang akan digunakan:
ollama pull llama3.24. 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,StartEventdipancarkan ke event loop. - Step
setup_querymenangkapStartEventdan memancarkanQueryEvent. - Step
analyze_sentimentdansummarize_textmendengarkan tipe event yang sama (QueryEvent). Kedua step ini dieksekusi secara terpisah dan paralel tanpa saling menunggu. - Step
aggregate_resultsbertindak sebagai pengumpul (Fan-In). Step ini dipanggil setiap kali salah satu dariSentimentResultEventatauSummaryResultEventdipancarkan. Step mengecek state padaContext; jika kedua data belum lengkap, step mengembalikanNoneuntuk tetap menunggu. Setelah seluruh data lengkap, step memancarkanStopEventyang 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:
- Tangkap exception pada step handler yang memanggil LLM.
- Evaluasi counter retry pada event payload.
- Jika retry belum melebihi batas maksimum (misal 3 kali), pancarkan
RetryEventbaru yang membawa informasi attempt counter dan delay waktu (Exponential Backoff). - Step khusus
handle_retrymenerimaRetryEvent, melakukanasyncio.sleep()sesuai delay, lalu memancarkan ulang event pemrosesan awal. - Jika batas retry terlampaui, alihkan ke
ErrorEventatau kembalikanStopEventberisi 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
Documentatau tensor embedding) ke dalam kelasEvent. 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.


