Event Streaming #
Di balik sistem-sistem besar yang kita gunakan sehari-hari — feed media sosial yang selalu terkini, deteksi penipuan kartu kredit yang terjadi dalam milidetik, dashboard analytics yang bergerak real-time — ada satu fondasi yang sama: aliran event yang terus mengalir dan diproses secara berkelanjutan. Event streaming bukan hanya cara untuk mengirim pesan dari satu service ke service lain. Ia adalah perubahan cara berpikir tentang data — dari sesuatu yang diam menunggu di-query, menjadi sesuatu yang bergerak dan terus mengalir. Artikel ini membahas event streaming dari prinsip dasarnya, perbedaan fundamentalnya dengan message queue dan database tradisional, anatomi platform streaming modern, empat use case nyata di industri, tantangan teknis yang sering diremehkan, hingga best practice yang membedakan implementasi yang bertahan dengan yang akhirnya menjadi masalah operasional.
Apa Itu Event Streaming? #
Event streaming adalah paradigma di mana event diproduksi, disimpan, dan dikonsumsi secara berkelanjutan oleh satu atau banyak sistem. Kata kunci yang membedakannya dari paradigma lain adalah berkelanjutan — stream tidak punya awal dan akhir yang terdefinisi seperti sebuah file atau batch query. Ia terus mengalir selama sistem berjalan.
Perbedaan cara berpikir yang paling fundamental adalah tentang apa yang disimpan:
Database tradisional — menyimpan state saat ini:
Tabel orders: {id: 1, status: "shipped", updated_at: "2026-04-17"}
✓ "Apa status order #1 sekarang?" → "shipped"
✗ "Kapan statusnya berubah dari paid ke shipped?" → tidak bisa dijawab
Event streaming — menyimpan urutan perubahan:
Stream order-events:
offset 0: OrderCreated {order_id: 1, total: 250000, at: "09:00"}
offset 1: PaymentReceived {order_id: 1, payment_id: "PAY-001", at: "09:05"}
offset 2: OrderPacked {order_id: 1, warehouse: "JKT-1", at: "10:30"}
offset 3: OrderShipped {order_id: 1, tracking: "JNE-123", at: "11:00"}
✓ "Apa status order #1 sekarang?" → replay semua event → "shipped"
✓ "Kapan transisi ke shipped?" → offset 3, pukul 11:00
✓ "Berapa lama dari payment ke shipped?" → 1 jam 55 menit
✓ "State order #1 pada pukul 10:00?" → replay sampai offset 1 → "paid"
Event streaming tidak fokus pada state akhir, melainkan pada perjalanan perubahan state itu sendiri. Ini yang membuka kemampuan seperti audit trail lengkap, time-travel debugging, dan rebuild state dari titik manapun dalam sejarah.
flowchart LR
subgraph DB["Database Tradisional"]
W1[Write] --> T1[("(Tabel:\nstate terkini)")]
T1 --> R1["Read\ncurrent state"]
end
subgraph ES["Event Streaming"]
W2[Write Event] --> L[("(Append-only\nEvent Log)")]
L --> R2["Read\ncurrent state\nreplay dari awal"]
L --> R3["Read\nhistorical state\nreplay sampai T"]
L --> R4["Read\nby multiple\nconsumer"]
endEvent Streaming vs Message Queue vs Database #
Ketiga paradigma ini sering dicampuradukkan. Memahami perbedaannya adalah syarat untuk memilih alat yang tepat — dan tidak over-engineering sistem sederhana dengan Kafka hanya karena terdengar canggih.
| Aspek | Database Tradisional | Message Queue | Event Streaming |
|---|---|---|---|
| Yang disimpan | State terkini | Task/pesan sementara | Urutan event permanen |
| Akses data | Pull (query) | Push ke consumer | Pull oleh consumer (offset) |
| Setelah dibaca | Tetap ada | Dihapus dari queue | Tetap ada (retention period) |
| Replay | Tidak applicable | Sangat sulit | Native — reset offset |
| Consumer | Banyak, baca bersamaan | Biasanya satu per message | Banyak, independen |
| Urutan | Tidak ada | Tergantung implementasi | Terjamin per partition |
| Skalabilitas | Vertikal | Horizontal dengan effort | Horizontal native |
| Cocok untuk | CRUD, transaksional | Task queue, background job | Pipeline, audit, analytics |
Analogi yang paling mudah untuk memahami perbedaan message queue vs event streaming:
- Message Queue seperti conveyor belt di pabrik — item diambil satu per satu oleh pekerja, dan setelah diambil hilang dari belt.
- Event Streaming seperti rekaman video — bisa ditonton berkali-kali oleh banyak penonton, penonton baru bisa mulai dari awal, dan penonton lama tidak terpengaruh.
flowchart TD
subgraph MQ["Message Queue — Consume & Delete"]
P1[Producer] --> Q["(Queue)"]
Q -->|"ambil msg1\n→ msg1 hilang"| C1[Consumer A]
Q -->|"ambil msg2\n→ msg2 hilang"| C2[Consumer B]
end
subgraph EV["Event Streaming — Log Persistent"]
P2[Producer] --> L[("(Event Log\noffset: 0,1,2,3...)")]
L -->|"baca dari offset 0\nevent tetap ada"| C3["Consumer A\nEmail Service"]
L -->|"baca dari offset 0\nindependen"| C4["Consumer B\nPayment Service"]
L -->|"baca dari offset 2\nbisa pilih mulai"| C5["Consumer C\nAnalytics"]
endAnatomi Platform Event Streaming #
Topic dan Partition #
Topic adalah named stream — channel logis tempat event sejenis dikumpulkan. order-events, user-events, payment-events adalah contoh topic.
Partition adalah unit fisik di dalam topic. Satu topic dibagi menjadi beberapa partition, masing-masing adalah log terurut yang independen. Partition adalah kunci skalabilitas horizontal.
flowchart LR
subgraph Topic["Topic: order-events (3 partitions)"]
subgraph P0["Partition 0"]
E0["offset 0\nOrderCreated\norder#1"] --> E3["offset 1\nOrderPaid\norder#1"] --> E6["offset 2\nOrderShipped\norder#1"]
end
subgraph P1["Partition 1"]
E1["offset 0\nOrderCreated\norder#2"] --> E4["offset 1\nOrderPaid\norder#2"] --> E7["offset 2\nOrderShipped\norder#2"]
end
subgraph P2["Partition 2"]
E2["offset 0\nOrderCreated\norder#3"] --> E5["offset 1\nOrderPaid\norder#3"]
end
end
PROD[Producer] -->|"key=order#1"| P0
PROD -->|"key=order#2"| P1
PROD -->|"key=order#3"| P2Pemilihan partition key adalah keputusan desain yang kritis — ia menentukan ordering dan distribusi beban:
// ANTI-PATTERN: tanpa partition key — event untuk order yang sama
// bisa masuk partition berbeda, ordering tidak terjamin
producer.SendMessage(&sarama.ProducerMessage{
Topic: "order-events",
Value: sarama.ByteEncoder(payload),
// tidak ada Key!
})
// BENAR: gunakan entity ID sebagai partition key
// Semua event untuk order yang sama selalu di partition yang sama
producer.SendMessage(&sarama.ProducerMessage{
Topic: "order-events",
Key: sarama.StringEncoder(event.OrderID), // ← ordering terjamin per order
Value: sarama.ByteEncoder(payload),
})
// Dampak: consumer selalu terima urutan yang benar:
// OrderCreated → PaymentReceived → OrderShipped
// bukan urutan acak karena masuk partition berbeda
Offset dan Consumer Group #
Offset adalah nomor urut setiap event dalam sebuah partition. Consumer mengontrol sendiri offset mana yang sudah dibaca — ini yang membuat replay menjadi native.
Consumer group adalah mekanisme horizontal scaling. Setiap consumer dalam grup mendapat subset partition yang berbeda — beban dibagi merata. Dua group berbeda bisa mengkonsumsi topic yang sama secara independen.
flowchart TD
subgraph T["Topic: order-events (6 partitions)"]
P0[P0]
P1[P1]
P2[P2]
P3[P3]
P4[P4]
P5[P5]
end
subgraph CG1["Consumer Group: payment-service (3 instances)"]
I1["Instance A\nP0, P1"]
I2["Instance B\nP2, P3"]
I3["Instance C\nP4, P5"]
end
subgraph CG2["Consumer Group: email-service (2 instances)"]
I4["Instance X\nP0, P1, P2"]
I5["Instance Y\nP3, P4, P5"]
end
P0 --> I1
P1 --> I1
P2 --> I2
P3 --> I2
P4 --> I3
P5 --> I3
P0 --> I4
P1 --> I4
P2 --> I4
P3 --> I5
P4 --> I5
P5 --> I5Jumlah instance consumer dalam satu group tidak bisa melebihi jumlah partition. Jika ada 6 partition dan 8 consumer instance, 2 instance akan idle. Desain jumlah partition harus mempertimbangkan kebutuhan scaling di masa depan — partition bisa ditambah tapi tidak bisa dikurangi tanpa rebalancing.
Retention dan Log Compaction #
Berbeda dengan message queue yang menghapus pesan setelah dikonsumsi, event streaming menyimpan event selama retention period yang dikonfigurasi.
| Strategi | Konfigurasi | Cocok Untuk |
|---|---|---|
| Time-based | retention.ms = 604800000 (7 hari) | Event transaksional, log |
| Size-based | retention.bytes = 10GB per partition | Storage terbatas |
| Log compaction | Simpan hanya event terbaru per key | State snapshots, config, profile |
| Infinite | retention.ms = -1 | Event sourcing, audit permanen |
// Log compaction: hanya event terbaru per key yang disimpan
// Cocok untuk: "current state" setiap entity
// Sebelum compaction:
// offset 0: UserUpdated {user_id: "123", name: "Umar"}
// offset 5: UserUpdated {user_id: "123", name: "Unis"}
// offset 9: UserUpdated {user_id: "123", name: "Unis Badri"}
//
// Setelah compaction:
// offset 9: UserUpdated {user_id: "123", name: "Unis Badri"}
// offset 0 dan 5 dihapus karena sudah ada yang lebih baru untuk key "123"
Empat Use Case Utama di Industri #
1. Microservices Communication #
Event streaming sebagai backbone komunikasi antar microservice — menggantikan synchronous API call untuk workflow yang tidak butuh respons langsung.
flowchart TD
subgraph Before["❌ Synchronous Chain — Rentan Cascading Failure"]
OS1[Order Service] -->|POST| PS1[Payment Service]
PS1 -->|POST| IS1[Inventory Service]
IS1 -->|POST| ES1[Email Service]
ES1 -- "down!" --> IS1
IS1 -- "error!" --> PS1
PS1 -- "error!" --> OS1
end
subgraph After["✅ Event Streaming — Fault Isolated"]
OS2[Order Service] -->|publish\norder.created| K["(Kafka)"]
K --> PS2[Payment Service]
K --> IS2[Inventory Service]
K --> ES2["Email Service\ndown? event antri\ndiproses saat pulih"]
end// ANTI-PATTERN: order service memanggil semua downstream langsung
func (s *OrderService) CreateOrder(req CreateOrderRequest) error {
order := s.db.CreateOrder(req)
s.paymentSvc.InitiatePayment(order.ID) // jika gagal → order gagal
s.inventorySvc.ReserveStock(order.Items) // jika gagal → order gagal
s.emailSvc.SendConfirmation(order.UserID) // jika gagal → order gagal
return nil
}
// BENAR: publish satu event, setiap service bereaksi independen
func (s *OrderService) CreateOrder(req CreateOrderRequest) error {
order := s.db.CreateOrder(req)
s.kafka.Publish("order.created", OrderCreatedEvent{
EventID: uuid.New().String(),
OrderID: order.ID,
UserID: order.UserID,
Items: order.Items,
CreatedAt: time.Now(),
})
return nil // selesai — tidak peduli siapa yang consume
}
2. Real-Time Analytics dan Stream Processing #
Event streaming memungkinkan analitik yang dihitung langsung dari stream, bukan dari batch query ke database.
// Konsep: hitung revenue per menit menggunakan stream processing
func buildRevenueStream(orderEvents KStream) KTable {
return orderEvents.
Filter(func(e OrderEvent) bool {
return e.Type == "order.paid"
}).
GroupBy(func(e OrderEvent) string {
return e.OccurredAt.Truncate(time.Minute).Format(time.RFC3339)
}).
Aggregate(
func() int64 { return 0 },
func(key string, event OrderEvent, agg int64) int64 {
return agg + event.Amount
},
)
// Hasil: tabel {minute → total_revenue} yang terupdate real-time
// tanpa menunggu batch job berjalan setiap jam
}
Use case nyata yang menggunakan pola ini: deteksi fraud real-time (pola transaksi anomali dalam window 5 menit), live leaderboard, monitoring latency API secara real-time, rekomendasi produk yang diperbarui saat user browsing.
3. Change Data Capture (CDC) #
CDC adalah teknik untuk mengambil setiap perubahan di database (INSERT, UPDATE, DELETE) dan mempublikasikannya sebagai event ke stream — tanpa polling atau trigger yang membebani database utama.
flowchart LR
DB[("(PostgreSQL\nDatabase)")] -->|"WAL / binlog\nreading"| D["Debezium\nConnector"]
D -->|"publish change events"| K["(Kafka)"]
K --> ES["Elasticsearch\nupdate search index"]
K --> RC["Redis\ninvalidate cache"]
K --> DW["Data Warehouse\nnew record"]
K --> AU["Audit Service\nlog perubahan"]// CDC event yang diproduksi Debezium untuk setiap UPDATE di database
type CDCEvent struct {
Before map[string]interface{} `json:"before"` // state sebelum
After map[string]interface{} `json:"after"` // state sesudah
Op string `json:"op"` // "c"=create, "u"=update, "d"=delete
Source CDCSource `json:"source"`
}
// Consumer memanfaatkan CDC untuk invalidate cache otomatis
func (c *CacheConsumer) HandleCDC(event CDCEvent) error {
if event.Op == "u" || event.Op == "d" {
userID := event.Before["id"].(string)
return c.redis.Del(ctx, "user:"+userID)
}
return nil
}
4. Event Sourcing sebagai Source of Truth #
Dalam event sourcing, stream event adalah database utama — bukan tambahan. State aplikasi direkonstruksi dengan me-replay event dari awal.
// ANTI-PATTERN: simpan hanya state akhir — kehilangan riwayat
func (s *AccountService) Debit(accountID string, amount int64) error {
account := s.db.FindAccount(accountID)
account.Balance -= amount
return s.db.Save(account)
// Pertanyaan "kenapa saldo jadi segini?" → tidak bisa dijawab
}
// BENAR: simpan event, rebuild state dari replay
func (s *AccountService) Debit(accountID string, amount int64) error {
event := AccountDebitedEvent{
AccountID: accountID,
Amount: amount,
OccurredAt: time.Now(),
}
return s.eventStore.Append(accountID, event)
}
func (s *AccountService) GetBalance(accountID string) int64 {
events := s.eventStore.LoadEvents(accountID)
account := Account{Balance: 0}
for _, e := range events {
switch ev := e.(type) {
case AccountCreditedEvent:
account.Balance += ev.Amount
case AccountDebitedEvent:
account.Balance -= ev.Amount
}
}
return account.Balance
// Bonus: bisa query "berapa saldo pada tanggal X?" dengan filter timestamp
}
Tantangan Teknis yang Sering Diremehkan #
Ordering Hanya Dijamin Per Partition #
Ini adalah salah satu kesalahpahaman paling umum. Event streaming tidak menjamin ordering global di seluruh topic. Ordering hanya dijamin di dalam satu partition.
// ANTI-PATTERN: tanpa partition key — event yang saling berhubungan
// bisa masuk partition berbeda, consumer lihat urutan yang salah
//
// Partition 0: OrderCreated(id=1), OrderShipped(id=2)
// Partition 1: PaymentReceived(id=1), OrderShipped(id=1)
//
// Consumer mungkin lihat: OrderShipped(id=1) sebelum OrderCreated(id=1)!
// BENAR: partition key = entity ID
//
// Partition 0: OrderCreated(id=1), PaymentReceived(id=1), OrderShipped(id=1)
// Partition 1: OrderCreated(id=2), PaymentReceived(id=2), OrderShipped(id=2)
//
// Semua event untuk order #1 selalu di Partition 0 → ordering terjamin
producer.SendMessage(&sarama.ProducerMessage{
Topic: "order-events",
Key: sarama.StringEncoder(event.OrderID), // ← kunci ordering
Value: sarama.ByteEncoder(payload),
})
Exactly-Once Semantics #
Exactly-once adalah jaminan yang paling sulit dan paling mahal. Dalam praktiknya, at-least-once + idempotent consumer adalah solusi yang jauh lebih realistis.
| Delivery Semantic | Jaminan | Trade-off | Kapan Digunakan |
|---|---|---|---|
| At-most-once | Mungkin hilang, tidak duplikat | Kehilangan data | Metrics non-kritis, log statistik |
| At-least-once | Pasti sampai, mungkin duplikat | Consumer harus idempotent | Paling umum — hampir semua use case |
| Exactly-once | Tidak hilang, tidak duplikat | Sangat mahal, overhead besar | Financial ledger, billing kritis |
// ANTI-PATTERN: commit offset sebelum processing — event bisa hilang
msg := consumer.Receive()
consumer.CommitOffset(msg.Offset) // ← commit terlalu cepat
processEvent(msg) // crash di sini = event hilang selamanya
// BENAR: commit offset SETELAH processing berhasil
func (c *Consumer) processLoop() {
for msg := range c.messages {
if err := c.processEvent(msg); err != nil {
// Jangan commit — message akan di-redeliver
log.Errorf("processing failed for offset %d: %v", msg.Offset, err)
continue
}
// Commit hanya setelah berhasil
c.session.MarkMessage(msg, "")
}
}
// Consumer idempotent — aman meski event datang 2x
func (c *Consumer) processEvent(msg *sarama.ConsumerMessage) error {
var event OrderEvent
json.Unmarshal(msg.Value, &event)
if processed, _ := c.idempotencyStore.Exists(ctx, event.EventID); processed {
return nil // sudah pernah diproses — skip dengan aman
}
if err := c.businessLogic(event); err != nil {
return err
}
c.idempotencyStore.Set(ctx, event.EventID, time.Now())
return nil
}
Schema Evolution dengan Schema Registry #
Ketika banyak producer dan consumer beroperasi di atas event yang sama, perubahan schema bisa merusak consumer yang belum di-deploy ulang.
flowchart TD
subgraph Without["❌ Tanpa Schema Registry"]
P1["Producer\nkirim format baru\namount_cents"] --> K1["(Kafka)"]
K1 --> C1["Consumer Lama\nharap amount\n💥 panic / error"]
end
subgraph With["✅ Dengan Schema Registry"]
P2["Producer\nregister schema baru"] --> SR[("(Schema\nRegistry)")]
SR -->|"cek kompatibilitas"| OK{Compatible?}
OK -->|Ya| K2["(Kafka)"]
OK -->|Tidak| REJ["❌ Reject\nbreaking change"]
K2 --> C2["Consumer\ndeserialize dengan\nschema yang benar"]
end// ANTI-PATTERN: rename field langsung — breaking change tanpa warning!
// Sebelum: {"user_id": "123", "amount": 250000}
// Sesudah: {"customer_id": "123", "amount_cents": 25000000} ← consumer lama crash!
// BENAR: backward-compatible evolution — tambah field baru, jangan hapus lama
// Schema v1
type OrderPaidEventV1 struct {
EventID string `json:"event_id"`
OrderID string `json:"order_id"`
UserID string `json:"user_id"` // tetap ada
Amount float64 `json:"amount"` // tetap ada
}
// Schema v2 — backward compatible
type OrderPaidEventV2 struct {
EventID string `json:"event_id"`
OrderID string `json:"order_id"`
UserID string `json:"user_id"` // tetap ada untuk compat
CustomerID string `json:"customer_id"` // field baru (alias)
Amount float64 `json:"amount"` // tetap ada untuk compat
AmountCents int64 `json:"amount_cents"` // field baru (alias)
PromoCode *string `json:"promo_code,omitempty"` // field baru, optional
}
// Consumer v1 tetap baca event v2 karena field lama masih ada
// Consumer v2 baca field baru yang lebih kaya
Jangan pernah mengubah atau menghapus field yang sudah ada di event schema yang sedang dikonsumsi production. Gunakan Confluent Schema Registry atau protokol seperti Protobuf / Avro yang memiliki aturan kompatibilitas built-in. Evolusi schema yang tidak terkontrol adalah salah satu penyebab insiden terbesar di sistem berbasis Kafka.
Best Practice Desain Event Streaming #
Event adalah Kontrak Publik #
Begitu sebuah event topic dipublikasikan dan ada consumer yang bergantung padanya, ia menjadi kontrak yang harus dijaga seperti public API.
Prinsip desain event sebagai kontrak publik:
□ Nama event berbasis fakta (kata kerja lampau): OrderCreated, bukan CreateOrder
□ Semua field wajib backward-compatible — tidak boleh dihapus atau ganti tipe
□ Field baru selalu optional dengan default value yang aman
□ Breaking change → buat event type baru dengan versi (order.created.v2)
□ Jalankan versi lama dan baru secara paralel selama masa transisi
□ Dokumentasikan di schema registry atau internal wiki
□ Jangan publish detail implementasi internal ke topic yang dikonsumsi service lain
Partition Key Strategy #
| Strategi | Key | Manfaat | Trade-off |
|---|---|---|---|
| Entity ID | order_id, user_id | Ordering per entity terjamin | Hot partition jika ada entity dengan traffic sangat tinggi |
| Hash-based | hash(entity_id) % n | Distribusi merata | Ordering tidak terjamin lintas partition |
| Round-robin | Tanpa key | Distribusi paling merata, throughput maksimal | Tidak ada ordering sama sekali |
| Time-based | date:YYYY-MM-DD | Mudah archiving per hari | Semua traffic hari ini ke satu partition |
Untuk sebagian besar use case bisnis, Entity ID adalah pilihan terbaik karena memberikan ordering semantik yang paling bermakna.
Consumer Lag Monitoring Wajib #
Consumer lag adalah metrik paling penting di event streaming — selisih antara offset terbaru yang diproduksi dengan offset terbaru yang sudah dikonsumsi.
Consumer lag = latest_offset - consumer_committed_offset
Lag = 0 → consumer real-time, tidak ada backlog ✓
Lag = 1.000 → 1000 event belum diproses
Lag terus naik → consumer tidak bisa mengimbangi producer
→ butuh scale out consumer atau optimize processing logic
// BENAR: monitor consumer lag dan alert saat melebihi threshold
func monitorConsumerLag(client sarama.Client, group, topic string) {
ticker := time.NewTicker(30 * time.Second)
for range ticker.C {
lag, err := calculateTotalLag(client, group, topic)
if err != nil {
log.Errorf("failed to calculate consumer lag: %v", err)
continue
}
// Kirim ke metrics system
metrics.Gauge("kafka.consumer.lag", float64(lag), map[string]string{
"group": group,
"topic": topic,
})
// Alert bertingkat
if lag > 100_000 {
alert.Critical(fmt.Sprintf("[KAFKA] consumer lag kritis: group=%s topic=%s lag=%d", group, topic, lag))
} else if lag > 10_000 {
alert.Warning(fmt.Sprintf("[KAFKA] consumer lag tinggi: group=%s topic=%s lag=%d", group, topic, lag))
}
}
}
Anti-Pattern yang Harus Dihindari #
// ✗ Event terlalu granular — overhead komunikasi tinggi, sulit dikelola
kafka.Publish("user.first-name.updated", event)
kafka.Publish("user.last-name.updated", event)
kafka.Publish("user.email.updated", event)
// ✓ Event di level bisnis yang bermakna
kafka.Publish("user.profile.updated", event) // satu event, satu intent
// ✗ Event berisi command, bukan fakta
kafka.Publish("send-welcome-email", event) // command!
kafka.Publish("process-payment", event) // command!
// ✓ Event berisi fakta historis
kafka.Publish("user.registered", event) // fakta
kafka.Publish("order.payment.received", event) // fakta
// ✗ Commit offset sebelum processing — event bisa hilang saat crash
session.MarkMessage(msg, "") // commit dulu
processEvent(msg) // crash di sini = event hilang
// ✓ Commit SETELAH processing berhasil
if err := processEvent(msg); err != nil { return err }
session.MarkMessage(msg, "") // commit setelah sukses
// ✗ Consumer tidak idempotent — event yang sama diproses 2x
func handle(event Event) { db.Insert(event.Data) } // duplikat jika redeliver!
// ✓ Cek event_id dulu
func handle(event Event) {
if idempotencyStore.Exists(event.EventID) { return nil }
db.Insert(event.Data)
idempotencyStore.Set(event.EventID)
}
// ✗ Tidak ada consumer lag monitoring — tahu masalah dari user complaint
// ✓ Alert otomatis saat lag > threshold, sebelum user merasakan dampak
// ✗ Schema berubah tanpa versioning — consumer lama tiba-tiba crash
// ✓ Schema Registry + backward-compatible evolution + versioned event type
Checklist Implementasi Event Streaming #
TOPIC DESIGN:
□ Nama topic mencerminkan domain, bukan implementasi (order-events, bukan kafka-order-svc)
□ Jumlah partition direncanakan dengan mempertimbangkan kebutuhan scaling
□ Retention policy sudah ditetapkan sesuai kebutuhan (7 hari, 30 hari, infinite)
□ Log compaction diaktifkan jika topic digunakan sebagai state store
PRODUCER:
□ Partition key ditetapkan untuk semua event yang butuh ordering
□ Event schema terdokumentasi di schema registry
□ Semua event punya: event_id, event_type, occurred_at, correlation_id
CONSUMER:
□ Commit offset SETELAH processing berhasil, bukan sebelum
□ Consumer idempotent — cek event_id sebelum proses
□ DLQ dikonfigurasi untuk event yang gagal melewati max retry
□ Consumer group name deskriptif (payment-service, bukan consumer-1)
SCHEMA:
□ Schema registry digunakan (Confluent Schema Registry / Protobuf)
□ Backward-compatible evolution untuk perubahan minor
□ Versioned event type untuk breaking change
□ Tidak ada penghapusan field dari schema yang sudah production
OBSERVABILITY:
□ Consumer lag monitoring aktif untuk semua consumer group
□ Alert consumer lag > threshold yang sudah ditetapkan
□ Correlation ID dipropagasi dari event ke log consumer
□ Distributed tracing mencakup alur dari producer ke consumer
Ringkasan #
- Event streaming menyimpan urutan perubahan, bukan state terkini — ini yang membuka audit trail lengkap, time-travel debugging, dan rebuild state dari titik manapun.
- Bukan message queue: event streaming menyimpan event setelah dikonsumsi (retention), mendukung banyak consumer independen, dan native replay; message queue menghapus pesan setelah dikonsumsi.
- Topic dan partition: topic adalah channel logis, partition adalah unit fisik yang memungkinkan horizontal scaling — jumlah partition adalah batas atas paralelisme consumer group.
- Gunakan entity ID sebagai partition key — ordering hanya dijamin per partition; tanpa key yang tepat, event yang saling berhubungan bisa diproses tidak berurutan.
- Consumer group: setiap group mendapat seluruh event secara independen; instance dalam satu group berbagi beban partition; instance tidak boleh melebihi jumlah partition.
- At-least-once + idempotent consumer adalah trade-off paling pragmatis — exactly-once sangat mahal dan jarang benar-benar diperlukan.
- Commit offset SETELAH processing — commit sebelum processing adalah resep kehilangan event saat consumer crash.
- Schema Registry wajib untuk lingkungan multi-team — mencegah breaking change schema yang tidak sengaja merusak consumer production.
- Consumer lag adalah metrik paling penting — monitoring dan alerting harus ada sebelum masalah terdeteksi oleh user, bukan sesudah.
- Event adalah kontrak publik — perlakukan seperti public API: backward-compatible untuk evolusi minor, versioning untuk breaking change.
← Sebelumnya: Event-Driven Berikutnya: Inversion of Control →