← Kembali ke Blog Kafka

Dari PHP Monolith ke Kafka Event-Driven: Cerita Nyata Migrasi di Crypto Exchange

20 Agustus 2026 • 12 min read • by Arif Kurniawan

Malam itu kami dapat alert: CPU server 98%, response time 30 detik. Itu bukan DDOS — itu adalah Monday effect dari crypto market yang bullish. Dalam 6 bulan, kami berhasil memotong p99 latency dari 8 detik jadi 200ms. Begini journey-nya.

Awal Mula: PHP Monolith yang Semakin Sesak

Boleh saya cerita sedikit tentang kondisi awal? Kami punya satu codebase PHP yang sudah berumur 7 tahun. Semua jadi satu: user management, wallet, trading engine, notification, reporting. Jutaan line code dalam folder yang sama tanpa boundaries yang jelas.

Saat market bullish, load tiba-tiba naik 10-20x dari kondisi normal. User login bersamaan, buka app, cek portofolio, deposit, withdraw — semuanya hitting satu database yang sama. Response time melambat drastis, timeout di mana-mana, dan yang paling menyakitkan: satu function yang error bisa menarik seluruh sistem ke bawah.

Kami sudah coba optimisasi standard — opcode cache, query optimization, horizontal scaling via load balancer. Tapi arsitektur monolitiknya itu sendiri yang jadi bottleneck. Adding more servers hanya menunda masalah, bukan solving it.

"Adding more servers to a monolith is like trying to unclog a pipe by adding more water pressure. At some point the pipe itself is the problem."

Design Decision: Kenapa Kafka?

Kami evaluasi beberapa pilihan: RabbitMQ, Redis Streams, langsung ke database, dan Kafka. Masing-masing punya trade-off.

RabbitMQ — bagus untuk message queuing standard, tapiKalau bicara soal crypto exchange yang butuh throughput tinggi dengan latency rendah, Redis Streams menarik karena sederhana dan satu binary dengan Redis yang sudah kami pakai buat caching. Tapi masalahnya: Redis Streams itu bounded queue, bukan distributed log. Retention dan replayability terbatas.

Kafka menang di beberapa aspek krusial untuk kami:

Kafka memungkinkan kami memecah monolith jadi event-driven services yang bisa scales independently. Satu service yang heavy load tidak akan ganggu service lain.

Topic Structure dan Consumer Groups

Design topic yang proper itu penting. Awalnya kami buat satu topik besar "events" — itu mistake. Setelah refactor, struktur topik kami jadi begini:

┌─────────────────────────────────────────────────────────────┐
│                     TOPIC STRUCTURE                        │
├─────────────────────────────────────────────────────────────┤
│  users.*          → User events (register, login, update)    │
│  wallets.*        → Wallet events (deposit, withdraw, tx)   │
│  trades.*         → Trading events (order, match, settle)    │
│  notifications.*  → Email, push, SMS events                  │
│  audit.*          → Audit log untuk compliance              │
└─────────────────────────────────────────────────────────────┘

Dengan naming convention seperti ini, governance lebih mudah. Partitioning strategy juga penting:

// Partitioning berdasarkan user_id untuk ensure ordering per user
// tapi tidak blocking antar users

func partitionKey(msg UserEvent) string {
    return msg.UserID  // semua events untuk user yang sama masuk partition yang sama
}

Consumer groups kami design agar satu service satu group ID. Ini memastikan setiap event diproses exactly once oleh setiap subscriber yang interested:

// Contoh consumer group configuration
type ConsumerConfig struct {
    Brokers:       []string{"kafka-1:9092", "kafka-2:9092", "kafka-3:9092"},
    GroupID:        "notification-service-v1",
    Topic:          "users.*",
    MinBytes:       10e3, // 10KB
    MaxBytes:       50e6, // 50MB
    MaxWait:        500 * time.Millisecond,
    OffsetReset:    sarama.OffsetNewest,
}

Producer dan Consumer Pattern di Go

Kami pakai segmentio/kafka-go sebagai client. Berikut pattern yang kami pakai untuk producer:

package kafka

import (
    "context"
    "time"
    "github.com/segmentio/kafka-go"
)

type Producer struct {
    writer *kafka.Writer
}

func NewProducer(brokers []string, topic string) *Producer {
    return &Producer{
        writer: &kafka.Writer{
            Addr:         kafka.TCP(brokers...),
            Topic:        topic,
            Balancer:     &kafka.LeastBytes{},
            BatchSize:    100,
            BatchTimeout: 10 * time.Millisecond,
            Async:        false, // sync untuk guarantee delivery
            RequiredAcks: kafka.RequireAll,
        },
    }
}

func (p *Producer) Publish(ctx context.Context, key, value []byte) error {
    return p.writer.WriteMessages(ctx, kafka.Message{
        Key:   key,
        Value: value,
        Time:  time.Now(),
    })
}

// Usage:
func (s *UserService) Register(ctx context.Context, user User) error {
    // 1. Insert to database first (source of truth)
    if err := s.repo.Create(ctx, user); err != nil {
        return err
    }

    // 2. Publish event ke Kafka (fire-and-forget dengan retry)
    event := UserRegisteredEvent{
        UserID:    user.ID,
        Email:     user.Email,
        Timestamp: time.Now(),
    }

    payload, _ := json.Marshal(event)
    if err := s.producer.Publish(ctx, []byte(user.ID), payload); err != nil {
        // Log error tapi jangan block — user sudah berhasil register
        s.logger.Error("failed to publish UserRegistered event", err)
    }

    return nil
}

Untuk consumer, kami pakai pola concurrent processing dengan worker pool:

type Consumer struct {
    reader   *kafka.Reader
    workers  int
    handler  MessageHandler
}

type MessageHandler func(ctx context.Context, msg kafka.Message) error

func (c *Consumer) Start(ctx context.Context) error {
    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        default:
            msg, err := c.reader.FetchMessage(ctx)
            if err != nil {
                continue
            }

            // Process dengan worker pool
            go func(m kafka.Message) {
                if err := c.handler(ctx, m); err != nil {
                    // Retry dengan exponential backoff
                    c.retry(m)
                }
                c.reader.CommitMessages(ctx, m)
            }(msg)
        }
    }
}

func (c *Consumer) retry(msg kafka.Message) {
    backoff := 1 * time.Second
    for i := 0; i < 3; i++ {
        time.Sleep(backoff)
        if err := c.handler(context.Background(), msg); err == nil {
            return
        }
        backoff *= 2
    }
    // Setelah 3 retry gagal, dead letter queue
    c.publishToDLQ(msg)
}

Challenge: Data Consistency

Inilah bagian tersulit dari migrasi. Bagaimana保证 consistency antara PHP monolith yang existing dengan services baru Go yang belum fully operational?

Masalah utama: dual-write. Kalau PHP menulis ke MySQL dan juga publish ke Kafka, apa yang terjadi kalau Kafka succeed tapi MySQL gagal? Atau sebaliknya?

Solusi yang kami pakai: Transactional Outbox Pattern.

// Tabel outbox untuk guarantee delivery
type Outbox struct {
    ID        string    `json:"id"`
    Aggregate string    `json:"aggregate"` // e.g., "user", "wallet"
    EventType string    `json:"event_type"` // e.g., "UserRegistered"
    Payload   []byte    `json:"payload"`
    Status    string    `json:"status"` // "pending", "published"
    CreatedAt time.Time `json:"created_at"`
    PublishedAt *time.Time `json:"published_at,omitempty"`
}

// Transactional Outbox — semua dalam satu transaction
func (s *UserService) Register(ctx context.Context, user User) error {
    return s.db.Transaction(ctx, func(tx *sql.Tx) error {
        // 1. Insert user
        if err := s.userRepo.WithTx(tx).Create(ctx, user); err != nil {
            return err
        }

        // 2. Insert outbox event (dalam transaction yang sama!)
        event := UserRegisteredEvent{...}
        payload, _ := json.Marshal(event)
        outbox := Outbox{
            ID:        uuid.New().String(),
            Aggregate: "user",
            EventType: "UserRegistered",
            Payload:   payload,
            Status:    "pending",
        }
        return s.outboxRepo.WithTx(tx).Create(ctx, outbox)
        // Transaction commits — kedua operasi atomic
    })
}

// Outbox Processor — baca dari outbox, publish ke Kafka, mark published
func (op *OutboxProcessor) Process(ctx context.Context) error {
    events, err := op.outboxRepo.GetPending(ctx, 100)
    if err != nil {
        return err
    }

    for _, event := range events {
        topic := fmt.Sprintf("%s.%s", event.Aggregate, event.EventType)
        if err := op.producer.Publish(ctx, []byte(event.ID), event.Payload); err != nil {
            continue // akan di-retry next cycle
        }

        event.Status = "published"
        event.PublishedAt = time.Now()
        op.outboxRepo.MarkPublished(ctx, event)
    }
    return nil
}

Dengan transactional outbox, kami guarantee: kalau database commit, event pasti ter-publish. Tidak ada lost update, tidak ada phantom events.

Rollout Strategy: Zero Downtime

Migrasi sekala ini tidak bisa langsung switch. Kami pakai Strangler Fig Pattern — pelan-pelan replace functionalities satu per satu.

Phase 1: Parallel Run
Services baru jalan di belakang, semua writes masih lewat PHP monolith. Kami capture events dari database change (via binlog) dan publish ke Kafka. Consumer services baca dari Kafka tapi tidak write ke primary database — hanya ke read replicas untuk testing.

Phase 2: Read Path Migration
Traffic reads mulai di-switch ke services baru. PHP masih handle writes, tapi reads dari mobile app, web dashboard dialihkan ke services Go. Ini berlangsung 2-3 minggu sampai confident.

Phase 3: Write Path Migration
Satu domain per minggu — mulai dari domain paling independent (notification, audit log) sampai yang paling critical (wallet, trading).

Phase 4: Decomission PHP
Semua writes sudah di services Go. PHP jadi historical, bisa di-keep untuk rollback kalau ada masalah. Setelah 30 hari clean, baru dimatikan.

Hasil dan Lessons Learned

Setelah 6 bulan migrasi:

Yang paling valuable: tim sekarang bisa work independently. Tim notification tidak perlu wait untuk tim wallet kalau mau deploy. Yang satu error, yang lain tidak jatuh.

"Kafka bukan magic bullet. Tapi dengan proper topic design, consumer patterns, dan rollout strategy, dia bisa jadi foundation untuk sistem yang scalable dan resilient."

Lessons learned yang worth share:

  1. Start with events, not services — design domain events dulu, baru decide service boundaries
  2. Transactional outbox itu bukan optional — kalau skip ini, data consistency akan hunt you down
  3. Monitoring harus dari day 1 — consumer lag, producer errors, dead letter queue size
  4. Incremental migration wins — big bang rewrite itu high risk. Strangler fig dengan feature flags jauh lebih aman

Migrasi ini bukan finish line — justru starting point. Sekarang dengan event-driven foundation, adding real-time features, cross-service analytics, dan ML pipelines jadi jauh lebih feasible.

Semoga cerita ini berguna untuk kamu yang lagi planning migrasi serupa. Kalau ada pertanyaan specific, boleh discuss di kolom komentar. 🚀