Replay Strategy #

Di distributed system, kegagalan bukan kemungkinan — melainkan kepastian yang tinggal menunggu waktu. Network akan timeout, service akan down, consumer akan crash di tengah proses, dan database sesekali kelebihan beban. Pertanyaan yang relevan bukan “apakah sistem kita akan gagal?” tapi “ketika gagal, apa yang terjadi pada data dan proses yang sedang berjalan?” Replay Strategy adalah jawaban atas pertanyaan itu: kemampuan sistem untuk memproses ulang event, message, atau request yang gagal — dengan aman, terkontrol, dan tanpa menimbulkan efek samping yang tidak diinginkan. Artikel ini membahas mengapa replay penting, lima jenis replay yang berbeda kegunaannya, tiga tantangan utama yang harus diantisipasi, dan best practice implementasi lengkap dengan kode konkret.

Apa Itu Replay Strategy? #

Replay Strategy adalah pendekatan desain sistem untuk memproses ulang event, message, atau request yang gagal, belum selesai, atau perlu diverifikasi ulang, dengan jaminan bahwa pengulangan tersebut tidak menyebabkan state yang rusak atau efek samping yang tidak dikehendaki.

Satu hal yang membedakan Replay Strategy dari sekadar “tambahkan retry”: replay adalah keputusan arsitektur, bukan beberapa baris kode for i < 3 { retry() }. Ia mencakup bagaimana event disimpan, bagaimana urutan dipertahankan, bagaimana duplikasi ditangani, dan bagaimana operator mengontrol proses replay di production.

Replay bisa terjadi di berbagai layer sistem:

LayerContoh ReplayTrigger Umum
Message QueueSQS, RabbitMQ retry saat consumer gagalConsumer crash, timeout
Event StreamingKafka consumer reset offset ke posisi lamaBug fix setelah deploy
HTTP / APIRetry request timeout ke downstreamNetwork instability
Background JobRerun scheduled job yang crash di tengahOOM, pod restart
CQRS CommandReplay command yang gagal dieksekusiValidation error upstream
Event SourcingRebuild seluruh state dari awal event logSchema migration, bug fix kalkulasi
flowchart TD
    A["Event / Request Masuk"] --> B{Berhasil Diproses?}
    B -- Ya --> C[Selesai ✓]
    B -- Tidak --> D{Tipe Kegagalan?}
    D -- Transient\ntimeout singkat --> E[Immediate Retry]
    D -- Butuh waktu\nrecover --> F[Exponential Backoff]
    D -- Max retry\nterlampaui --> G[Dead Letter Queue]
    D -- Bug fix\ndeploy --> H["Manual Replay\nEvent Store"]
    D -- Rebuild state --> I["Event Sourcing\nReplay"]
    E --> B
    F --> B
    G --> J[Alert + Review]
    J --> H

Mengapa Replay Strategy Itu Kritis #

Ada tiga alasan fundamental yang membuat replay bukan fitur opsional di sistem modern.

Kegagalan transient adalah norma, bukan pengecualian. Downstream service yang sedang overload, koneksi database yang sesaat terputus, rate limit dari third-party API — semua ini adalah kegagalan sementara yang seharusnya bisa di-recover secara otomatis. Tanpa replay, setiap kegagalan transient berpotensi menjadi kehilangan data permanen.

Replay adalah fondasi dari at-least-once delivery. Hampir semua message broker modern — Kafka, SQS, RabbitMQ, Pub/Sub — menjamin at-least-once delivery, bukan exactly-once. Artinya sistem harus siap menerima event yang sama lebih dari sekali, dan harus bisa “mengulang” proses setelah kegagalan. Tanpa replay strategy yang benar, sistem hanya bekerja di kondisi ideal.

Replay adalah alat operasional yang sangat berharga. Ini yang sering diabaikan: replay bukan hanya untuk error recovery. Engineer menggunakan replay untuk memproses ulang event setelah bug fix di-deploy, untuk rebuild read model setelah migrasi schema, untuk re-sync data setelah insiden, dan untuk audit ulang transaksi yang dicurigai.

sequenceDiagram
    participant Prod as Event Producer
    participant Queue as Message Queue
    participant Consumer
    participant DB

    Prod->>Queue: Publish event (order.created)
    Queue->>Consumer: Deliver event
    Consumer->>DB: Proses — DB timeout!
    DB--xConsumer: Error
    Consumer-->>Queue: NACK — tidak ack
    Queue->>Consumer: Redeliver (retry ke-1)
    Consumer->>DB: Proses lagi
    DB-->>Consumer: OK
    Consumer-->>Queue: ACK ✓
    Note over Prod,DB: Tanpa replay: event hilang.\nDengan replay: event tetap aman.

Perbedaan mendasar antara sistem yang mature dan yang belum: sistem yang mature mengasumsikan kegagalan akan terjadi dan mendesain replay sejak awal. Sistem yang belum mature baru menambahkan replay setelah insiden pertama — dan sering terlambat karena event-nya sudah tidak tersimpan.


Jenis-Jenis Replay Strategy #

Tidak semua replay sama. Setiap jenis punya use case, kelebihan, dan risikonya sendiri. Memilih jenis yang tepat tergantung pada konteks kegagalan dan karakteristik sistem.

1. Immediate Retry #

Replay langsung dilakukan setelah kegagalan, tanpa delay. Cocok untuk kegagalan yang benar-benar transient dan singkat — koneksi TCP yang sesaat putus, lock contention yang segera terlepas.

// ANTI-PATTERN: retry tanpa batas, tanpa delay — membombardir downstream
func processEvent(event Event) error {
    for {
        err := handler.Process(event)
        if err == nil {
            return nil
        }
        // jika downstream down, ini loop selamanya dan memperparah kondisi
    }
}

// BENAR: retry dengan batas maksimum yang jelas
func processWithRetry(event Event, maxAttempts int) error {
    var lastErr error
    for attempt := 1; attempt <= maxAttempts; attempt++ {
        lastErr = handler.Process(event)
        if lastErr == nil {
            return nil
        }
        log.Warnf("attempt %d/%d failed for event %s: %v",
            attempt, maxAttempts, event.ID, lastErr)
    }
    return fmt.Errorf("all %d attempts failed: %w", maxAttempts, lastErr)
}
Immediate retry hanya aman untuk kegagalan yang sembuh dalam hitungan milidetik. Jangan gunakan immediate retry jika downstream membutuhkan lebih dari beberapa detik untuk recover — gunakan exponential backoff.

2. Delayed Retry dengan Exponential Backoff #

Replay dilakukan setelah delay yang semakin lama setiap kali gagal. Ini memberi downstream service waktu untuk recover tanpa dibombardir oleh retry storm.

// ANTI-PATTERN: fixed delay — tidak adaptif terhadap lama kegagalan
func processWithFixedRetry(event Event) error {
    for i := 0; i < 5; i++ {
        if err := handler.Process(event); err == nil {
            return nil
        }
        time.Sleep(1 * time.Second) // selalu 1 detik, tidak peduli berapa lama down
    }
    return errors.New("max retries exceeded")
}

// BENAR: exponential backoff dengan jitter mencegah thundering herd
func processWithBackoff(ctx context.Context, event Event) error {
    baseDelay := 1 * time.Second
    maxDelay  := 5 * time.Minute
    maxAttempts := 8

    for attempt := 1; attempt <= maxAttempts; attempt++ {
        err := handler.Process(event)
        if err == nil {
            return nil
        }

        if attempt == maxAttempts {
            return fmt.Errorf("exhausted %d attempts: %w", maxAttempts, err)
        }

        // Delay: 1s, 2s, 4s, 8s, 16s, 32s, 60s (capped)...
        delay := baseDelay * time.Duration(1<<uint(attempt-1))
        if delay > maxDelay {
            delay = maxDelay
        }

        // Jitter ±25%: agar tidak semua consumer retry pada waktu yang persis sama
        jitter := time.Duration(rand.Int63n(int64(delay / 4)))
        if rand.Intn(2) == 0 {
            delay += jitter
        } else {
            delay -= jitter
        }

        log.Infof("retry attempt %d in %v for event %s", attempt+1, delay, event.ID)

        select {
        case <-ctx.Done():
            return ctx.Err()
        case <-time.After(delay):
        }
    }
    return nil
}

Pola delay dengan backoff:

AttemptBase DelayDengan Jitter (±25%)
1 → 21 detik0.75 – 1.25 detik
2 → 32 detik1.5 – 2.5 detik
3 → 44 detik3 – 5 detik
4 → 58 detik6 – 10 detik
5 → 616 detik12 – 20 detik
6 → 732 detik24 – 40 detik
7 → 860 detik (capped)45 – 75 detik

3. Dead Letter Queue (DLQ) #

Ketika event sudah melebihi batas retry dan tetap gagal, ia dipindahkan ke Dead Letter Queue — antrian terpisah untuk event-event bermasalah yang tidak bisa diproses saat ini. DLQ adalah safety net yang memastikan event tidak hilang.

flowchart LR
    A[Queue Utama] --> B[Consumer]
    B -- Gagal --> C[Retry 1]
    C -- Gagal --> D[Retry 2]
    D -- Gagal --> E["Retry 3\nMax Reached"]
    E --> F[Dead Letter Queue]
    F --> G["Alert / Monitor"]
    G --> H{"Root Cause\nDitemukan?"}
    H -- Ya, fix deployed --> I["Manual Replay\ndari DLQ"]
    H -- Perlu investigasi --> J["Engineer Review\nEvent Detail"]
    I --> A
// ANTI-PATTERN: buang event jika gagal — kehilangan data
func (c *Consumer) handleMessage(msg Message) error {
    if err := c.process(msg); err != nil {
        log.Errorf("failed to process %s, dropping: %v", msg.ID, err)
        return nil // event hilang selamanya
    }
    return nil
}

// BENAR: kirim ke DLQ dengan metadata lengkap, jangan buang event
func (c *Consumer) handleWithDLQ(ctx context.Context, msg Message) error {
    err := c.processWithBackoff(ctx, msg)
    if err == nil {
        return nil
    }

    // Enrich dengan metadata untuk debugging nanti
    dlqMsg := DLQMessage{
        OriginalMessage: msg,
        FailureReason:   err.Error(),
        AttemptCount:    msg.ApproximateReceiveCount,
        FailedAt:        time.Now(),
        ConsumerVersion: c.version,
        StackTrace:      debug.Stack(),
    }

    if dlqErr := c.dlqQueue.Send(ctx, dlqMsg); dlqErr != nil {
        // DLQ send gagal juga — log CRITICAL, butuh perhatian segera
        log.Criticalf("FAILED to send to DLQ, event %s will be lost: %v",
            msg.ID, dlqErr)
        return dlqErr
    }

    log.Warnf("event %s sent to DLQ after %d attempts: %v",
        msg.ID, msg.ApproximateReceiveCount, err)
    return nil // ack message utama agar tidak looping di queue utama
}

4. Manual Replay dari Event Store #

Event disimpan secara durable dan bisa diputar ulang secara selektif kapan saja — berdasarkan timestamp, event ID, topic/partition, atau kriteria bisnis tertentu.

// ANTI-PATTERN: replay tanpa throttle — membanjiri downstream sekaligus
func replayAllFromDLQ() {
    events := dlq.GetAll() // bisa jutaan event
    for _, e := range events {
        process(e) // semua diproses secepat mungkin — sistem bisa down
    }
}

// BENAR: replay terkontrol dengan rate limiter dan progress tracking
func (r *ReplayService) ReplayFromDLQ(ctx context.Context, opts ReplayOptions) error {
    log.Infof("starting replay: from=%v to=%v filter=%s rate=%f/s dryRun=%v",
        opts.From, opts.To, opts.Filter, opts.RatePerSecond, opts.DryRun)

    events, err := r.dlqStore.Query(ctx, opts.From, opts.To, opts.Filter)
    if err != nil {
        return fmt.Errorf("query DLQ: %w", err)
    }

    log.Infof("found %d events to replay", len(events))

    limiter := rate.NewLimiter(rate.Limit(opts.RatePerSecond), 1)
    var succeeded, failed, skipped int

    for i, event := range events {
        // Bisa dibatalkan kapan saja lewat context
        if err := limiter.Wait(ctx); err != nil {
            log.Infof("replay cancelled at event %d/%d", i, len(events))
            return err
        }

        if opts.DryRun {
            log.Infof("[DRY RUN] would replay event %s", event.ID)
            skipped++
            continue
        }

        if err := r.processor.Process(ctx, event); err != nil {
            log.Errorf("replay failed for event %s: %v", event.ID, err)
            failed++
            continue
        }
        succeeded++

        // Progress log setiap 100 event
        if succeeded%100 == 0 {
            log.Infof("replay progress: %d/%d (success=%d fail=%d skip=%d)",
                i+1, len(events), succeeded, failed, skipped)
        }
    }

    log.Infof("replay complete: total=%d succeeded=%d failed=%d skipped=%d",
        len(events), succeeded, failed, skipped)
    return nil
}

5. Event Sourcing Replay #

Dalam arsitektur event sourcing, event log adalah source of truth — state aplikasi direkonstruksi ulang dari urutan event sejak awal. Replay di sini berarti membangun ulang projection atau read model dari awal.

// ANTI-PATTERN: rebuild langsung di tabel production tanpa backup
func rebuildProjection() {
    db.Exec("DELETE FROM order_summary") // bahaya! tidak ada rollback jika gagal
    for _, event := range eventLog.GetAll() {
        applyEvent(event)
    }
}

// BENAR: rebuild ke shadow table, swap atomis setelah selesai
func (p *ProjectionRebuilder) RebuildOrderProjection(ctx context.Context) error {
    shadowTable := fmt.Sprintf("order_summary_rebuild_%d", time.Now().Unix())

    // Buat shadow table terlebih dahulu
    if err := p.db.CreateTableLike(shadowTable, "order_summary"); err != nil {
        return fmt.Errorf("create shadow table: %w", err)
    }

    // Stream semua event dari awal ke shadow table
    events, err := p.eventLog.StreamAll(ctx, "orders", StreamOptions{
        FromOffset: 0,
        BatchSize:  500,
    })
    if err != nil {
        p.db.DropTable(shadowTable) // cleanup
        return fmt.Errorf("stream events: %w", err)
    }

    var count int
    for event := range events {
        if err := p.applyEventToTable(ctx, shadowTable, event); err != nil {
            p.db.DropTable(shadowTable)
            return fmt.Errorf("apply event %d: %w", event.Sequence, err)
        }
        count++
        if count%1000 == 0 {
            log.Infof("rebuilt %d events...", count)
        }
    }

    // Atomic swap: rename shadow → production
    if err := p.db.RenameTable("order_summary", "order_summary_old"); err != nil {
        return err
    }
    if err := p.db.RenameTable(shadowTable, "order_summary"); err != nil {
        return err
    }
    p.db.DropTable("order_summary_old")

    log.Infof("projection rebuild complete: %d events applied", count)
    return nil
}

Tantangan Utama dalam Replay #

Tiga tantangan ini harus diantisipasi saat merancang replay strategy — mengabaikannya adalah sumber bug yang sulit dilacak.

Duplicate Processing #

Replay hampir selalu menghasilkan duplikasi — event yang sama diproses lebih dari sekali. Ini bukan bug dari replay-nya, tapi konsekuensi yang harus ditangani dengan idempotency.

// ANTI-PATTERN: consumer yang tidak idempotent — replay = bencana
func (c *Consumer) handlePayment(event PaymentRequestedEvent) error {
    // Jika event ini di-replay, payment diproses dua kali → double charge
    return c.paymentGateway.Charge(event.UserID, event.Amount)
}

// BENAR: consumer idempotent — cek event_id sebelum proses
func (c *Consumer) handlePayment(ctx context.Context, event PaymentRequestedEvent) error {
    processed, err := c.processedEvents.Exists(ctx, event.EventID)
    if err != nil {
        return fmt.Errorf("check processed: %w", err)
    }
    if processed {
        log.Infof("event %s already processed, skipping", event.EventID)
        return nil // aman di-skip saat replay
    }

    // Proses dalam transaction: charge + tandai sebagai processed — atomic
    return c.db.Transaction(func(tx *gorm.DB) error {
        if err := c.paymentGateway.Charge(event.UserID, event.Amount); err != nil {
            return err
        }
        return tx.Create(&ProcessedEvent{
            EventID:     event.EventID,
            ProcessedAt: time.Now(),
        }).Error
    })
}
flowchart TD
    A[Event Datang untuk Diproses] --> B{"event_id ada\ndi processed_events?"}
    B -- Ya --> C["Skip — return nil ✓\nReplay aman"]
    B -- Tidak --> D["Proses Event\ndalam DB Transaction"]
    D --> E[Eksekusi operasi bisnis]
    E --> F["Insert event_id\nke processed_events"]
    F --> G{"Transaction\nCommit?"}
    G -- Ya --> H[ACK ke queue ✓]
    G -- Tidak --> I["Rollback\nEvent akan di-redeliver"]
    I --> A

Ordering dan Urutan Event #

Replay bisa mengubah urutan pemrosesan event. Ini kritis untuk event yang saling bergantung — misalnya OrderCreated harus selalu diproses sebelum OrderShipped.

// ANTI-PATTERN: publish tanpa partition key — urutan tidak terjamin lintas partition
func (p *Producer) publishOrderEvent(event OrderEvent) error {
    msg := &sarama.ProducerMessage{
        Topic: "order-events",
        // Tanpa key: Kafka distribusi ke partition secara round-robin
        // OrderCreated (partition 0) bisa diproses setelah OrderShipped (partition 1)
        Value: sarama.ByteEncoder(mustMarshal(event)),
    }
    _, _, err := p.client.SendMessage(msg)
    return err
}

// BENAR: gunakan entity ID sebagai partition key
// Semua event untuk order yang sama dijamin masuk partition yang sama
func (p *Producer) publishOrderEvent(event OrderEvent) error {
    msg := &sarama.ProducerMessage{
        Topic: "order-events",
        Key:   sarama.StringEncoder(event.OrderID), // ← kunci ordering
        Value: sarama.ByteEncoder(mustMarshal(event)),
    }
    _, _, err := p.client.SendMessage(msg)
    return err
}

Side Effects yang Tidak Bisa Di-undo #

Ini tantangan paling licin: beberapa side effect tidak idempotent secara alami — mengirim email, memotong saldo, memanggil webhook ke partner. Replay yang tidak didesain dengan baik akan memicu semua ini berulang kali.

// ANTI-PATTERN: side effect langsung di handler tanpa cek state
func (h *OrderHandler) handleOrderConfirmed(event OrderConfirmedEvent) error {
    h.updateDatabase(event)
    h.emailService.SendConfirmation(event.UserEmail) // dikirim ulang setiap replay!
    h.webhookService.Notify(event)                   // dipanggil ulang setiap replay!
    return nil
}

// BENAR: cek state dulu — side effect hanya dipanggil jika state belum berubah
func (h *OrderHandler) handleOrderConfirmed(ctx context.Context, event OrderConfirmedEvent) error {
    order, err := h.orderRepo.FindByID(event.OrderID)
    if err != nil {
        return err
    }

    if order.Status == "CONFIRMED" {
        // State sudah benar — semua side effect sudah dijalankan sebelumnya
        // Replay aman: tidak ada email duplikat, tidak ada webhook duplikat
        return nil
    }

    // Update state terlebih dahulu
    if err := h.orderRepo.UpdateStatus(event.OrderID, "CONFIRMED"); err != nil {
        return err
    }

    // Side effects hanya dieksekusi sekali karena state sudah berubah
    h.emailService.SendConfirmation(event.UserEmail)
    h.webhookService.Notify(event)
    return nil
}

Best Practice Implementasi #

Simpan Event Secara Durable #

Replay hanya mungkin dilakukan jika event masih ada. Banyak sistem yang tidak menyimpan event dengan benar dan baru menyadarinya saat insiden.

Persyaratan minimal event store yang mendukung replay:

PersyaratanAlasan
Persistent storage, bukan hanya in-memoryEvent tidak hilang saat restart
Event tidak dihapus setelah diprosesReplay butuh event asli
Metadata lengkap: ID, timestamp, versi schema, sourceDebugging dan filtering saat replay
Bisa di-query berdasarkan timestamp/ID/kriteria bisnisReplay selektif, bukan harus semua event
Retention policy jelasStorage tidak membengkak tanpa batas

Throttle Replay Massal #

Replay massal tanpa kontrol adalah cara tercepat untuk menjatuhkan production system sendiri.

// ANTI-PATTERN: replay semua DLQ sekaligus — flood sistem
func replayAll() {
    events := dlq.GetAll() // bisa jutaan event
    for _, e := range events {
        process(e) // makan semua resource downstream sekaligus
    }
}

// BENAR: throttle + dry run + kill switch lewat context
type ReplayConfig struct {
    RatePerSecond float64       // berapa event per detik diproses
    BatchSize     int           // berapa event per batch
    DryRun        bool          // simulasi tanpa eksekusi nyata
    From          time.Time     // filter dari waktu tertentu
    To            time.Time     // filter sampai waktu tertentu
    Filter        string        // filter berdasarkan event type atau criteria lain
}

func (r *ReplayService) ReplayWithControl(ctx context.Context, cfg ReplayConfig) error {
    if cfg.DryRun {
        log.Info("[DRY RUN] no events will actually be processed")
    }
    limiter := rate.NewLimiter(rate.Limit(cfg.RatePerSecond), cfg.BatchSize)

    for _, event := range r.fetchEvents(cfg) {
        // Kill switch: ctx.Cancel() dari luar menghentikan replay kapan saja
        if err := limiter.Wait(ctx); err != nil {
            log.Info("replay stopped by cancellation")
            return err
        }
        if !cfg.DryRun {
            r.processor.Process(ctx, event)
        }
    }
    return nil
}

Observability adalah Wajib #

Replay tanpa visibility adalah operasi buta. Kamu tidak tahu apakah berhasil, berapa yang sudah diproses, atau apakah ada yang gagal lagi.

// BENAR: log dan metrics yang informatif untuk setiap replay operation
type ReplayMetrics struct {
    TotalEvents  int
    Succeeded    int
    Failed       int
    Skipped      int           // idempotency skip
    Duration     time.Duration
    TriggeredBy  string        // siapa yang trigger: "alerting-system", "engineer-unis"
    ReplayReason string        // mengapa: "bug-fix-deploy-v2.3.1", "db-migration"
}

func (r *ReplayService) logCompletion(m ReplayMetrics) {
    log.Infof("[REPLAY COMPLETE] reason=%q triggered_by=%q "+
        "total=%d succeeded=%d failed=%d skipped=%d duration=%v",
        m.ReplayReason, m.TriggeredBy,
        m.TotalEvents, m.Succeeded, m.Failed, m.Skipped, m.Duration)

    // Kirim ke metrics system untuk dashboard dan alerting
    metrics.Gauge("replay.success_rate",
        float64(m.Succeeded)/float64(m.TotalEvents)*100)
    metrics.Histogram("replay.duration_seconds", m.Duration.Seconds())

    // Alert jika failure rate terlalu tinggi
    failureRate := float64(m.Failed) / float64(m.TotalEvents)
    if failureRate > 0.05 { // > 5% gagal
        alerting.Send(fmt.Sprintf("High replay failure rate: %.1f%%", failureRate*100))
    }
}

Kapan Menggunakan Jenis Replay yang Mana #

flowchart TD
    A[Butuh Replay] --> B{"Kegagalan\ntransient singkat?"}
    B -- Ya, < 1 detik --> C["Immediate Retry\nmaks 3x"]
    B -- Tidak --> D{"Downstream butuh\nwaktu recover?"}
    D -- Ya --> E["Exponential Backoff\ndengan jitter"]
    E --> F{"Max retry\nterlampaui?"}
    F -- Ya --> G[Dead Letter Queue]
    D -- Tidak --> H{"Bug sudah\ndiperbaiki?"}
    H -- Ya, butuh proses\nulang event historis --> I["Manual Replay\ndari Event Store"]
    H -- Butuh rebuild\nseluruh state --> J["Event Sourcing\nReplay"]
    G --> K[Alert + Engineer Review]
    K --> I
SituasiReplay yang Tepat
Kegagalan transient, sembuh dalam milidetikImmediate retry (≤ 3x)
Downstream perlu waktu recovery (overload, deployment)Exponential backoff + jitter
Event gagal terus setelah max retryDead Letter Queue
Bug fix di-deploy, perlu proses ulang event historisManual replay dari event store
Migrasi schema, rebuild read model dari nolEvent Sourcing replay
Rate limit dari third-party APIDelayed retry dengan backoff
Insiden besar, ribuan event di DLQThrottled batch replay dengan monitoring

Anti-Pattern yang Harus Dihindari #

// ✗ Replay tanpa idempotency — menghasilkan double processing
for _, e := range dlqEvents {
    processPayment(e) // tidak ada pengecekan apakah sudah diproses!
}
// ✓ Selalu cek event_id sebelum memproses ulang

// ✗ DLQ tanpa metadata — susah debugging saat investigasi
dlq.Send(Message{Body: originalMsg.Body})
// ✓ Sertakan failure reason, attempt count, dan timestamp lengkap
dlq.Send(DLQMessage{
    Body:          originalMsg.Body,
    FailureReason: err.Error(),
    AttemptCount:  originalMsg.ReceiveCount,
    FailedAt:      time.Now(),
    StackTrace:    debug.Stack(),
})

// ✗ Replay massal tanpa kill switch — tidak bisa dihentikan darurat
func replayAllDLQ() {
    for _, e := range getAllEvents() {
        process(e) // tidak ada cara menghentikan jika sistem mulai kewalahan
    }
}
// ✓ Selalu gunakan context sebagai kill switch
func replayAllDLQ(ctx context.Context) {
    for _, e := range getAllEvents() {
        if ctx.Err() != nil {
            log.Info("replay cancelled")
            return
        }
        process(ctx, e)
    }
}

// ✗ Event disimpan hanya di in-memory queue — hilang saat restart
queue := make(chan Event, 1000) // tidak persistent
// ✓ Event store yang durable: PostgreSQL, S3, atau Kafka dengan retention panjang

// ✗ Side effect tanpa guard — email dikirim setiap replay
func handleOrderConfirmed(event Event) {
    emailService.Send(event.UserEmail) // kirim setiap kali dipanggil
}
// ✓ Cek state dulu, side effect hanya jika state belum berubah

Checklist Implementasi Replay Strategy #

DESAIN EVENT STORE:
  □ Event disimpan di persistent storage dengan retention yang jelas
  □ Setiap event punya ID unik, timestamp, versi schema, dan source
  □ Event bisa di-query berdasarkan waktu, ID, tipe, atau kriteria bisnis
  □ Event tidak dihapus setelah diproses — hanya ditandai

CONSUMER IDEMPOTENCY:
  □ Semua consumer mengecek event_id sebelum memproses
  □ Operasi bisnis dan pencatatan event_id dalam satu DB transaction
  □ Consumer tetap mengembalikan sukses jika event sudah diproses (skip)

DEAD LETTER QUEUE:
  □ DLQ dikonfigurasi untuk semua queue produksi
  □ Setiap event di DLQ menyimpan: failure reason, attempt count, timestamp
  □ Alert terpasang ketika DLQ menerima event baru
  □ Proses review DLQ terdokumentasi di runbook

MANUAL REPLAY:
  □ Tools replay tersedia dan terdokumentasi
  □ Replay mendukung dry run mode
  □ Replay punya rate limiter dan bisa dibatalkan (context cancel)
  □ Progress dan hasil replay di-log dengan lengkap

SIDE EFFECTS:
  □ Email, webhook, dan charge diproteksi dengan state check
  □ Side effect hanya dieksekusi jika state belum berubah
  □ Test replay dijalankan di staging sebelum production

OBSERVABILITY:
  □ Setiap replay operation punya log: triggered_by, reason, result
  □ Metrics: success rate, failure count, duration
  □ Alert jika failure rate replay > threshold

Ringkasan #

  • Replay Strategy adalah keputusan arsitektur — bukan sekadar retry loop; mencakup cara menyimpan event, menjaga urutan, menangani duplikasi, dan mengontrol proses replay di production.
  • Kegagalan adalah kepastian — desain replay sejak awal, jangan tunggu insiden pertama.
  • Lima jenis replay — immediate retry (transient singkat), exponential backoff (perlu waktu recover), DLQ (gagal permanen), manual replay (bug fix/audit), event sourcing replay (rebuild state).
  • Idempotency adalah syarat mutlak replay — tanpa idempotency, replay adalah resep double processing dan data rusak.
  • DLQ adalah safety net wajib — event tidak boleh hilang hanya karena belum bisa diproses saat ini; sertakan metadata lengkap untuk debugging.
  • Exponential backoff + jitter mencegah retry storm — jangan retry langsung saat downstream butuh waktu recovery; jitter mencegah thundering herd.
  • Throttle replay massal — replay ribuan event tanpa rate limiting bisa menjatuhkan production; selalu sediakan kill switch lewat context.
  • Side effects butuh guard — cek state terlebih dahulu sebelum mengirim email, webhook, atau charge; side effect hanya boleh terjadi sekali.
  • Observability wajib — log triggered_by, reason, success/failure count, dan duration untuk setiap replay operation.

← Sebelumnya: Race Condition   Berikutnya: Reactive Programming →

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