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:
| Layer | Contoh Replay | Trigger Umum |
|---|---|---|
| Message Queue | SQS, RabbitMQ retry saat consumer gagal | Consumer crash, timeout |
| Event Streaming | Kafka consumer reset offset ke posisi lama | Bug fix setelah deploy |
| HTTP / API | Retry request timeout ke downstream | Network instability |
| Background Job | Rerun scheduled job yang crash di tengah | OOM, pod restart |
| CQRS Command | Replay command yang gagal dieksekusi | Validation error upstream |
| Event Sourcing | Rebuild seluruh state dari awal event log | Schema 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 --> HMengapa 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:
| Attempt | Base Delay | Dengan Jitter (±25%) |
|---|---|---|
| 1 → 2 | 1 detik | 0.75 – 1.25 detik |
| 2 → 3 | 2 detik | 1.5 – 2.5 detik |
| 3 → 4 | 4 detik | 3 – 5 detik |
| 4 → 5 | 8 detik | 6 – 10 detik |
| 5 → 6 | 16 detik | 12 – 20 detik |
| 6 → 7 | 32 detik | 24 – 40 detik |
| 7 → 8 | 60 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 --> AOrdering 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:
| Persyaratan | Alasan |
|---|---|
| Persistent storage, bukan hanya in-memory | Event tidak hilang saat restart |
| Event tidak dihapus setelah diproses | Replay butuh event asli |
| Metadata lengkap: ID, timestamp, versi schema, source | Debugging dan filtering saat replay |
| Bisa di-query berdasarkan timestamp/ID/kriteria bisnis | Replay selektif, bukan harus semua event |
| Retention policy jelas | Storage 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| Situasi | Replay yang Tepat |
|---|---|
| Kegagalan transient, sembuh dalam milidetik | Immediate retry (≤ 3x) |
| Downstream perlu waktu recovery (overload, deployment) | Exponential backoff + jitter |
| Event gagal terus setelah max retry | Dead Letter Queue |
| Bug fix di-deploy, perlu proses ulang event historis | Manual replay dari event store |
| Migrasi schema, rebuild read model dari nol | Event Sourcing replay |
| Rate limit dari third-party API | Delayed retry dengan backoff |
| Insiden besar, ribuan event di DLQ | Throttled 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 →