Temporal dengan Python: Durable Execution untuk Workflow yang Andal
Temporal adalah platform untuk durable execution: ia memungkinkan Anda menulis logika bisnis yang berjalan lama dan bersifat stateful sebagai kode biasa yang tetap bertahan ketika proses crash, saat deployment, maupun ketika infrastruktur gagal. Alih-alih merangkai queue, cron job, dan database untuk melacak "sampai mana proses order ini berjalan," Anda cukup menulis sebuah fungsi dan Temporal menjamin fungsi itu berjalan sampai selesai persis seperti yang ditulis, bahkan jika mesin yang menjalankannya mati di tengah jalan.
Apa Itu Durable Execution
Sebagian besar tool workflow yang mungkin pernah Anda lihat — Airflow, Prefect, Dagster — adalah scheduler untuk pipeline data. Mereka sangat baik dalam menjalankan DAG berisi task secara berkala dan menampilkan hasilnya. Celery adalah task queue untuk memindahkan pekerjaan latar belakang. Temporal menyelesaikan masalah yang berbeda: menjaga satu proses berumur panjang tetap benar di tengah berbagai kegagalan.
Ide intinya adalah model event-sourced replay:
- Sebuah Workflow adalah logika bisnis Anda yang ditulis sebagai kode deterministik. Temporal tidak menyimpan variabel workflow di memori selamanya. Sebaliknya, setiap langkah penting (hasil activity, timer yang berbunyi, signal yang masuk) ditambahkan ke event history yang durable.
- Jika proses worker crash, Temporal memulai workflow kembali pada worker lain dan mereplay event history tersebut. Setiap baris kode Anda dieksekusi ulang, tetapi alih-alih memanggil activity lagi, SDK mengembalikan hasil yang sudah tercatat. Ketika replay menyusul ke titik tempat crash terjadi, eksekusi berlanjut normal.
- Karena replay, kode workflow harus deterministik: dengan history yang sama, ia harus mengambil jalur yang sama setiap kali. Apa pun yang non-deterministik (panggilan jaringan, angka acak, membaca jam, query database) harus terjadi di dalam sebuah Activity.
Sebuah Activity adalah fungsi biasa untuk side effect dan pekerjaan non-deterministik. Activity tidak direplay; hasilnya dicatat sekali dan dipakai ulang. Activity adalah tempat kode Anda berbicara dengan dunia luar.
Pembagian inilah yang membuat mekanismenya bekerja: workflow adalah "otak" yang durable dan bisa direplay, sedangkan activity adalah "tangan" yang menyentuh sistem eksternal.
Arsitektur Inti
Sebuah deployment Temporal memiliki beberapa komponen:
- Temporal Server (Cluster): backend yang menyimpan event history, menjadwalkan task, serta menegakkan timeout dan retry. Ia adalah sumber kebenaran, didukung oleh database (PostgreSQL, MySQL, atau Cassandra).
- Task Queue: queue bernama yang dipakai server untuk menyerahkan pekerjaan ke kode Anda. Worker mem-poll task queue; klien dan server merutekan task workflow dan activity ke queue tersebut.
- Worker: proses yang Anda jalankan untuk menampung kode workflow dan activity Anda. Worker mem-poll task queue, mengeksekusi task, dan melaporkan hasil kembali ke server. Kode Anda ada di sini, bukan di server.
- Client: cara kode aplikasi memulai workflow, mengirim signal, dan membaca state.
- Web UI: dashboard untuk memeriksa workflow yang berjalan dan yang selesai, event history-nya, input, output, dan kegagalannya.
Server tidak pernah menjalankan kode Anda. Ia hanya mengorkestrasi. Pemisahan ini berarti Anda dapat men-deploy versi worker baru, menskalakan worker secara horizontal, dan server tetap menjaga history tetap aman.
Menyiapkan Server Pengembangan
Untuk pengembangan lokal, Temporal CLI menyertakan dev server mandiri dengan database in-memory dan Web UI.
# Instal Temporal CLI (macOS / Linux)
curl -sSf https://temporal.download/cli.sh | sh
Atau dengan Homebrew
brew install temporal
Jalankan dev server lokal dengan Web UI di http://localhost:8233
temporal server start-dev
Dev server mendengarkan koneksi SDK di localhost:7233 dan menyajikan Web UI di localhost:8233. Ia me-reset state saat restart, yang wajar untuk pengembangan. Untuk produksi Anda menjalankan cluster sungguhan atau menggunakan Temporal Cloud (dibahas nanti).
Instal Python SDK:
pip install temporalio
Contoh yang Konsisten: Order Fulfillment
Sepanjang tutorial ini kita membangun satu workflow: memproses order pelanggan. Langkahnya adalah menagih pembayaran, mereservasi inventory, mengirim paket, dan memberi tahu pelanggan. Setiap langkah dapat gagal dan sebaiknya di-retry. Nanti kita tambahkan durable timer, signal untuk membatalkan, dan query untuk mengecek status.
Mendefinisikan Activity
Activity menampung semua side effect. Masing-masing adalah fungsi biasa (boleh async) yang didekorasi dengan @activity.defn.
# activities.py
import asyncio
from dataclasses import dataclass
from temporalio import activity
@dataclass
class OrderInput:
orderid: str
customeremail: str
amountcents: int
@activity.defn
async def chargepayment(order: OrderInput) -> str:
activity.logger.info(f"Menagih {order.amountcents} untuk {order.orderid}")
# Kode nyata akan memanggil payment gateway di sini.
await asyncio.sleep(0.2)
return f"charge{order.orderid}"
@activity.defn
async def reserveinventory(orderid: str) -> bool:
await asyncio.sleep(0.2)
return True
@activity.defn
async def shippackage(orderid: str) -> str:
await asyncio.sleep(0.5)
return f"tracking{orderid}"
@activity.defn
async def sendnotification(email: str, message: str) -> None:
activity.logger.info(f"Email ke {email}: {message}")
Activity bisa gagal dan Temporal akan me-retry sesuai policy. Karena hasilnya tercatat secara durable, workflow yang di-retry atau direplay tidak pernah menagih kartu yang sama dua kali untuk langkah yang sama yang sudah tercatat.
Mendefinisikan Workflow
Workflow mengorkestrasi activity. Ia didekorasi dengan @workflow.defn, dan entry point-nya dengan @workflow.run.
# workflows.py
from datetime import timedelta
from temporalio import workflow
from temporalio.common import RetryPolicy
Impor stub activity melalui pass-through yang aman untuk sandbox.
with workflow.unsafe.importspassedthrough():
from activities import (
OrderInput,
chargepayment,
reserveinventory,
shippackage,
sendnotification,
)
@workflow.defn
class OrderWorkflow:
def init(self) -> None:
self.status = "started"
@workflow.run
async def run(self, order: OrderInput) -> str:
retry = RetryPolicy(
initialinterval=timedelta(seconds=1),
backoffcoefficient=2.0,
maximuminterval=timedelta(seconds=30),
maximumattempts=5,
)
self.status = "charging"
chargeid = await workflow.executeactivity(
chargepayment,
order,
starttoclosetimeout=timedelta(seconds=10),
retrypolicy=retry,
)
self.status = "reserving"
await workflow.executeactivity(
reserveinventory,
order.orderid,
starttoclosetimeout=timedelta(seconds=10),
retrypolicy=retry,
)
self.status = "shipping"
tracking = await workflow.executeactivity(
shippackage,
order.orderid,
starttoclosetimeout=timedelta(seconds=30),
retrypolicy=retry,
)
await workflow.executeactivity(
sendnotification,
args=[order.customeremail, f"Order dikirim: {tracking}"],
starttoclosetimeout=timedelta(seconds=10),
retrypolicy=retry,
)
self.status = "completed"
return tracking
Perhatikan workflow tidak pernah memanggil requests, time.time(), atau random secara langsung. Ia hanya meng-await activity dan primitif SDK. Itulah penerapan nyata dari kendala determinisme.
Kendala Determinisme
Ketika Anda membutuhkan sesuatu yang non-deterministik di dalam kode workflow, gunakan pengganti deterministik dari SDK, yang mencatat nilainya ke history:
# Di dalam workflow — benar
now = workflow.now() # bukan datetime.now()
requestid = workflow.uuid4() # bukan uuid.uuid4()
delay = workflow.random().randint(1, 5) # bukan random.randint
Di dalam workflow — SALAH, merusak replay
import datetime, uuid, random
now = datetime.datetime.now() # nilai berbeda saat replay -> error non-determinisme
Python SDK menjalankan kode workflow dalam sandbox yang memblokir banyak kesalahan semacam ini, tetapi Anda tetap harus memperlakukan aturan ini sebagai hal mendasar: tidak ada I/O, tidak ada jam, tidak ada randomness, tidak ada global mutable state di kode workflow. Letakkan semuanya di activity.
Menjalankan Worker
Worker terhubung ke server, mendaftarkan workflow dan activity Anda, lalu mem-poll task queue.
# worker.py
import asyncio
from temporalio.client import Client
from temporalio.worker import Worker
from workflows import OrderWorkflow
from activities import (
chargepayment,
reserveinventory,
shippackage,
sendnotification,
)
async def main() -> None:
client = await Client.connect("localhost:7233")
worker = Worker(
client,
taskqueue="order-fulfillment",
workflows=[OrderWorkflow],
activities=[
chargepayment,
reserveinventory,
shippackage,
sendnotification,
],
)
await worker.run()
if name == "main":
asyncio.run(main())
Jalankan dengan python worker.py. Worker tetap hidup mem-poll queue order-fulfillment. Menskalakan throughput cukup dengan menjalankan lebih banyak proses worker yang menunjuk ke queue yang sama.
Memulai Workflow dari Client
Kode aplikasi (handler API, sebuah skrip) memulai workflow melalui client.
# startorder.py
import asyncio
from temporalio.client import Client
from workflows import OrderWorkflow
from activities import OrderInput
async def main() -> None:
client = await Client.connect("localhost:7233")
order = OrderInput(
order
id="A-1001",
customeremail="customer@example.com",
amountcents=4999,
)
# executeworkflow memulai dan menunggu hasilnya.
result = await client.executeworkflow(
OrderWorkflow.run,
order,
id="order-A-1001", # workflow ID
taskqueue="order-fulfillment",
)
print("Tracking:", result)
if name == "main":
asyncio.run(main())
Workflow ID dan Idempotensi
id adalah pengenal unik dari sebuah eksekusi workflow. Ia menjadi idempotency key Anda. Jika Anda mencoba memulai workflow dengan ID yang sedang berjalan, secara default Temporal menolak duplikatnya. Artinya Anda dapat dengan aman me-retry "mulai order A-1001" dari pemanggil yang tidak andal tanpa pernah menciptakan dua proses fulfillment untuk order yang sama. Gunakan pengenal bisnis (order ID) alih-alih nilai acak.
Untuk memulai tanpa menunggu, gunakan startworkflow dan dapatkan handle-nya:
handle = await client.startworkflow(
OrderWorkflow.run,
order,
id="order-A-1001",
task
queue="order-fulfillment",
)
result = await handle.result() # await nanti jika diperlukan
Retry dan Timeout
Temporal membedakan beberapa jenis timeout. Yang paling sering Anda atur adalah starttoclosetimeout: waktu maksimum satu percobaan activity boleh berjalan sebelum dianggap gagal dan di-retry. RetryPolicy mengatur bagaimana kegagalan di-retry, dengan exponential backoff.
from temporalio.common import RetryPolicy
from datetime import timedelta
retry = RetryPolicy(
initialinterval=timedelta(seconds=1),
backoffcoefficient=2.0,
maximuminterval=timedelta(minutes=1),
maximumattempts=10, # 0 berarti tak terbatas
nonretryableerrortypes=["InvalidCardError"],
)
Secara default activity di-retry tanpa batas dengan backoff, sehingga gangguan sementara pada payment gateway akan teratasi dengan sendirinya begitu gateway pulih; workflow menunggu dengan sabar. Atur maximumattempts untuk membatasinya.
Heartbeat untuk Activity yang Berjalan Lama
Untuk activity yang berjalan lama (memproses file besar, memanggil model yang lambat), gunakan heartbeat agar server tahu activity masih hidup dan dapat mendeteksi worker yang macet dengan cepat. Pasangkan heartbeat dengan heartbeattimeout.
@activity.defn
async def processlargebatch(items: list[str]) -> int:
processed = 0
for i, item in enumerate(items):
# lakukan pekerjaan...
processed += 1
activity.heartbeat(i) # laporkan progres; mengaktifkan deteksi gagal cepat
return processed
await workflow.executeactivity(
process
largebatch,
items,
start
toclosetimeout=timedelta(minutes=30),
heartbeattimeout=timedelta(seconds=30),
)
Jika worker crash, heartbeat yang terlewat membuat server menjadwalkan ulang activity dengan cepat alih-alih menunggu sampai starttoclosetimeout penuh. Saat retry, activity dapat membaca activity.info().heartbeatdetails untuk melanjutkan dari titik terakhir.
Durable Timer
Sebuah workflow dapat tidur selama apa pun — detik, hari, atau bulan — tanpa menahan resource apa pun. workflow.sleep membuat durable timer yang disimpan di history. Worker boleh mati selama tidur; ketika timer berbunyi, Temporal membangunkan workflow pada worker mana pun yang tersedia.
@workflow.run
async def run(self, order: OrderInput) -> str:
# ... kirim paket ...
# Tunggu 7 hari, lalu minta review. Tidak ada proses yang harus tetap berjalan.
await workflow.sleep(timedelta(days=7))
await workflow.executeactivity(
sendnotification,
args=[order.customeremail, "Bagaimana order Anda?"],
starttoclosetimeout=timedelta(seconds=10),
)
return tracking
Ini berbeda secara mendasar dari cron. Cron berbunyi mengikuti jam dan Anda harus mencari state di database untuk tahu apa yang harus dilakukan. Durable timer adalah bagian dari satu workflow berkesinambungan yang sudah memegang seluruh konteksnya dalam variabel lokal. Tidak ada state eksternal yang perlu direkonsiliasi.
Signal, Query, dan Update
Workflow yang sedang berjalan bukan kotak hitam. Anda dapat berinteraksi dengannya.
Sebuah Signal mengirim data ke workflow yang berjalan secara asinkron. Sifatnya fire-and-forget dan dapat mengubah state workflow. Di sini kita izinkan pelanggan membatalkan sebelum pengiriman.
@workflow.defn
class OrderWorkflow:
def init(self) -> None:
self.status = "started"
self.cancelled = False
@workflow.signal
def cancelorder(self) -> None:
self.cancelled = True
@workflow.query
def status(self) -> str:
return self.status
@workflow.run
async def run(self, order: OrderInput) -> str:
self.status = "charging"
await workflow.executeactivity(
chargepayment, order,
starttoclosetimeout=timedelta(seconds=10),
)
# Tunggu hingga satu jam untuk kemungkinan pembatalan sebelum pengiriman.
try:
await workflow.waitcondition(
lambda: self.cancelled,
timeout=timedelta(hours=1),
)
except TimeoutError:
pass
if self.cancelled:
self.status = "cancelled"
await workflow.executeactivity(
sendnotification,
args=[order.customeremail, "Order Anda dibatalkan."],
starttoclosetimeout=timedelta(seconds=10),
)
return "cancelled"
self.status = "shipping"
tracking = await workflow.executeactivity(
shippackage, order.orderid,
starttoclosetimeout=timedelta(seconds=30),
)
self.status = "completed"
return tracking
Sebuah Query membaca state workflow secara sinkron tanpa mengubahnya. Query tidak boleh memutasi state atau memanggil activity. Sebuah Update adalah interaksi request/response yang lebih baru dan tervalidasi, yang dapat mengubah state sekaligus mengembalikan hasil, dengan validasi opsional sebelum diterima.
Dari sisi client:
handle = client.getworkflowhandle("order-A-1001")
Query status saat ini
current = await handle.query(OrderWorkflow.status)
print(current)
Kirim signal pembatalan
await handle.signal(OrderWorkflow.cancelorder)
Child Workflow dan continueasnew
Sebuah workflow dapat memulai child workflow untuk memecah masalah besar atau memberi sub-proses ID, retry policy, dan siklus hidupnya sendiri.
tracking = await workflow.executechildworkflow(
ShippingWorkflow.run,
order.order
id,
id=f"shipping-{order.orderid}",
)
Event history bertambah pada setiap langkah. Workflow yang berulang selamanya (langganan yang menagih bulanan, agent yang berjalan banyak iterasi) pada akhirnya akan menumpuk history yang sangat besar, sehingga memperlambat replay. Solusinya adalah continueasnew: ia mengakhiri eksekusi saat ini dan secara atomik memulai eksekusi baru dengan ID yang sama dan history yang bersih, hanya membawa state yang Anda berikan.
@workflow.run
async def run(self, state: int) -> None:
for in range(1000):
await workflow.executeactivity(
doperiodicwork,
starttoclosetimeout=timedelta(seconds=30),
)
await workflow.sleep(timedelta(days=30))
state += 1
# Pangkas history dan lanjutkan dengan state yang dibawa.
workflow.continueasnew(state)
Penanganan Error
Exception dari activity muncul di workflow sebagai ActivityError yang membungkus ApplicationError. Untuk menandai error sebagai permanen agar Temporal tidak me-retry-nya, lempar ApplicationError yang non-retryable dari activity.
from temporalio.exceptions import ApplicationError
@activity.defn
async def chargepayment(order: OrderInput) -> str:
if order.amountcents <= 0:
# Request yang buruk tidak akan pernah berhasil saat retry.
raise ApplicationError("Jumlah tidak valid", nonretryable=True)
...
Di workflow Anda dapat menangkap dan melakukan kompensasi, yang merupakan cara membangun sebuah saga (membatalkan langkah sebelumnya ketika langkah berikutnya gagal):
from temporalio.exceptions import ActivityError
try:
await workflow.executeactivity(reserveinventory, order.orderid, ...)
except ActivityError:
# Kompensasi: refund tagihan sebelumnya.
await workflow.executeactivity(refundpayment, chargeid, ...)
raise
Versioning dan Patching
Karena workflow mereplay history lama, mengubah kode workflow dapat memecah eksekusi yang sedang berjalan: kode baru mungkin mengambil jalur yang berbeda dari yang diharapkan history yang sudah tercatat. Untuk perubahan yang aman, gunakan workflow.patched, yang mencatat cabang mana yang diambil sebuah eksekusi.
if workflow.patched("use-express-shipping"):
tracking = await workflow.executeactivity(expressship, ..., )
else:
tracking = await workflow.executeactivity(shippackage, ..., )
Eksekusi lama mereplay cabang asli; eksekusi baru mengambil cabang baru. Setelah semua workflow lama selesai, Anda memanggil workflow.deprecatepatch("use-express-shipping") dan kemudian menghapus kondisionalnya. Bagi banyak tim, menjalankan worker dengan Worker Versioning (mengunci build ID) adalah strategi pelengkap.
Pengujian dengan Time-Skipping Environment
Temporal menyertakan test environment yang menjalankan server in-memory dan dapat melompati waktu, sehingga workflow dengan timer tujuh hari selesai dalam milidetik di test suite Anda.
# testorder.py
import uuid
import pytest
from temporalio.testing import WorkflowEnvironment
from temporalio.worker import Worker
from workflows import OrderWorkflow
from activities import (
OrderInput, charge
payment, reserveinventory,
ship
package, sendnotification,
)
@pytest.mark.asyncio
async def test
ordercompletes():
async with await WorkflowEnvironment.start
timeskipping() as env:
async with Worker(
env.client,
task
queue="test-queue",
workflows=[OrderWorkflow],
activities=[chargepayment, reserveinventory,
shippackage, sendnotification],
):
order = OrderInput("A-1", "a@b.com", 4999)
result = await env.client.executeworkflow(
OrderWorkflow.run,
order,
id=f"test-{uuid.uuid4()}",
taskqueue="test-queue",
)
assert result.startswith("tracking")
Anda juga dapat me-mock activity dengan mendaftarkan fungsi pengganti pada test worker, sehingga dapat menguji logika workflow tanpa menyentuh sistem eksternal sungguhan.
Catatan tentang AI dan Workflow Agent
Durable execution cocok dengan pipeline AI modern. Sebuah agent LLM multi-langkah — mengambil konteks, memanggil model, menjalankan tool, mengevaluasi, mungkin mengulang, lalu memanggil model lagi — persis jenis proses berjalan lama, rawan gagal, dan stateful yang menjadi tujuan Temporal. Panggilan model dan invokasi tool menjadi activity dengan retry dan timeout-nya sendiri, sehingga error rate-limit atau tool yang labil tidak menghilangkan seluruh run. Logika orkestrasi (langkah mana berikutnya, kapan berhenti, berapa banyak iterasi) berada dalam workflow deterministik, dan continueasnew menjaga history loop agent yang panjang tetap terbatas. Persetujuan human-in-the-loop menjadi signal yang ditunggu workflow dengan durable timer sebagai deadline.
Bentuk yang sama mencakup saga bisnis: sebuah order, pengajuan pinjaman, atau alur onboarding yang membentang berhari-hari, memanggil banyak layanan, dan harus melakukan kompensasi dengan rapi saat gagal.
Temporal Cloud vs Self-Hosting
Anda dapat menjalankan Temporal Server sendiri di Kubernetes atau VM, didukung PostgreSQL atau Cassandra Anda sendiri. Ini memberi kendali penuh tetapi berarti Anda mengoperasikan cluster, database-nya, skalanya, dan upgrade-nya. Temporal Cloud adalah penawaran terkelola: Anda hanya menjalankan worker dan menghubungkannya ke namespace yang di-host melalui mTLS, sementara Temporal mengoperasikan server, storage, dan skalanya. Jalur umum adalah memulai dari dev server lokal, memvalidasi di cluster self-hosted kecil atau namespace Cloud, lalu memilih berdasarkan apakah mengoperasikan backend sepadan dengan kendali yang diberikannya.
Praktik Terbaik dan Kesalahan Umum
- Jaga workflow tetap deterministik. Tidak ada akses jam, randomness, jaringan, file, atau database secara langsung di kode workflow. Gunakan
workflow.now,workflow.uuid4,workflow.random, dan activity. Pelanggaran determinisme adalah bug produksi yang paling umum. - Letakkan semua side effect di activity. Jika menyentuh dunia luar atau bisa berbeda antar run, tempatnya di activity.
- Gunakan workflow ID yang bermakna. Kaitkan ID dengan entitas bisnis untuk idempotensi alami dan pencarian mudah di Web UI.
- Atur timeout secara eksplisit. Selalu atur
starttoclosetimeoutpada activity; tambahkanheartbeattimeoutuntuk yang berjalan lama dan kirim heartbeat dari dalamnya. - Tandai kegagalan permanen sebagai non-retryable. Input yang malformed atau respons 4xx sebaiknya melempar
nonretryable=Trueagar tidak di-retry selamanya. - Batasi history panjang dengan continueasnew. Setiap loop atau workflow berumur panjang sebaiknya secara berkala melakukan continue-as-new.
- Versikan perubahan berisiko dengan patching. Jangan pernah mengubah urutan activity pada workflow yang sudah ter-deploy tanpa
workflow.patchedselama eksekusi lama masih berjalan. - Jangan memblokir thread workflow. Hindari
time.sleepdan panggilan blocking sinkron; gunakanworkflow.sleepdan activity async. - Uji dengan time skipping. Cakup jalur yang banyak timer dan berbasis signal di test environment alih-alih menunggu di waktu nyata.
Kesimpulan
Temporal mengubah proses multi-langkah yang rapuh menjadi kode biasa yang sederhananya tidak kehilangan progres. Dengan memisahkan logika menjadi workflow deterministik dan activity yang ber-side-effect, serta merekonstruksi state melalui replay event history, ia menghilangkan rangka tambahan berupa queue, tabel status, dan job rekonsiliasi yang biasa.
Poin-poin kunci:
- Durable execution berarti kode workflow Anda bertahan dari crash dan restart melalui event-sourced replay.
- Workflow harus deterministik; semua I/O dan non-determinisme masuk ke activity.
- Server mengorkestrasi dan menyimpan history; worker menjalankan kode Anda; task queue menghubungkannya.
- Retry, timeout, heartbeat, dan durable timer sudah bawaan, bukan tambahan.
- Signal, query, dan update memungkinkan Anda berinteraksi dengan workflow yang hidup; child workflow dan
continueasnewmenjaga proses besar tetap terkelola. - Ia jauh lebih cocok untuk saga bisnis berjalan lama dan pipeline AI/agent multi-langkah dibanding scheduler pipeline data atau task queue biasa.