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"]
    end

Event 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.

AspekDatabase TradisionalMessage QueueEvent Streaming
Yang disimpanState terkiniTask/pesan sementaraUrutan event permanen
Akses dataPull (query)Push ke consumerPull oleh consumer (offset)
Setelah dibacaTetap adaDihapus dari queueTetap ada (retention period)
ReplayTidak applicableSangat sulitNative — reset offset
ConsumerBanyak, baca bersamaanBiasanya satu per messageBanyak, independen
UrutanTidak adaTergantung implementasiTerjamin per partition
SkalabilitasVertikalHorizontal dengan effortHorizontal native
Cocok untukCRUD, transaksionalTask queue, background jobPipeline, 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"]
    end

Anatomi 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"| P2

Pemilihan 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 --> I5
Jumlah 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.

StrategiKonfigurasiCocok Untuk
Time-basedretention.ms = 604800000 (7 hari)Event transaksional, log
Size-basedretention.bytes = 10GB per partitionStorage terbatas
Log compactionSimpan hanya event terbaru per keyState snapshots, config, profile
Infiniteretention.ms = -1Event 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 SemanticJaminanTrade-offKapan Digunakan
At-most-onceMungkin hilang, tidak duplikatKehilangan dataMetrics non-kritis, log statistik
At-least-oncePasti sampai, mungkin duplikatConsumer harus idempotentPaling umum — hampir semua use case
Exactly-onceTidak hilang, tidak duplikatSangat mahal, overhead besarFinancial 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 #

StrategiKeyManfaatTrade-off
Entity IDorder_id, user_idOrdering per entity terjaminHot partition jika ada entity dengan traffic sangat tinggi
Hash-basedhash(entity_id) % nDistribusi merataOrdering tidak terjamin lintas partition
Round-robinTanpa keyDistribusi paling merata, throughput maksimalTidak ada ordering sama sekali
Time-baseddate:YYYY-MM-DDMudah archiving per hariSemua 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 →

About | Author | Content Scope | Editorial Policy | Privacy Policy | Disclaimer | Contact