argabayu— Backend & Systems Engineer

ingest-pipeline

Pipeline ingest event yang tahan terhadap lonjakan dan tidak pernah menggandakan data.
Event/hari
40M
p99 latency
82ms
Duplikat
0

Pipeline ingest event untuk produk analitik internal. Menangani sekitar 40 juta event per hari dari empat belas layanan berbeda.

Masalah

Pipeline lama menulis langsung ke PostgreSQL dari tiap producer. Setiap ada lonjakan traffic, connection pool habis dan seluruh penulisan — termasuk yang tidak terkait — ikut gagal. Tidak ada deduplikasi, jadi retry klien menciptakan baris ganda.

Keputusan

Saya memilih menambahkan lapisan antrian di depan basis data, bukan memperbesar pool koneksi. Memperbesar pool hanya menggeser batasnya; antrian mengubah bentuk beban dari spike menjadi aliran rata.

Yang saya tolak:

  • Kafka — butuh operasi tambahan yang tidak sanggup dirawat tim berdua.
  • Menulis batch besar tiap beberapa menit — menambah latensi yang tidak perlu untuk kasus ini.
  • Memakai ORM — query-nya harus bisa saya baca saat insiden jam 3 pagi.

Yang dipakai: antrean dalam proses berbasis channel dengan flush berkala, plus kunci idempotensi per event_id.

go
// Flush mengosongkan buffer ke basis data dalam satu transaksi.
// Kunci idempotensi berada di ON CONFLICT, bukan di lapisan aplikasi —
// jadi dua proses yang menulis event sama tidak akan menggandakan baris.
func (w *Writer) Flush(ctx context.Context) error {
	tx, err := w.db.BeginTx(ctx, nil)
	if err != nil {
		return err
	}
	defer tx.Rollback()

	stmt, err := tx.PrepareContext(ctx, `
		INSERT INTO events (id, kind, payload, received_at)
		VALUES ($1, $2, $3, $4)
		ON CONFLICT (id) DO NOTHING`)
	if err != nil {
		return err
	}
	defer stmt.Close()

	for _, e := range w.buf {
		if _, err := stmt.ExecContext(ctx, e.ID, e.Kind, e.Payload, e.At); err != nil {
			return fmt.Errorf("flush %s: %w", e.ID, err)
		}
	}
	w.buf = w.buf[:0]
	return tx.Commit()
}

Kunci ON CONFLICT (id) DO NOTHING adalah bagian terpenting: idempotensi dijamin oleh basis data, bukan oleh hati-hatinya kode aplikasi.

Hasil

  • Lonjakan tidak lagi menggagalkan penulisan layanan lain.
  • p99 turun dari sekitar 1,4 detik ke 82ms.
  • Nol duplikat sejak diluncurkan — diverifikasi dengan query harian yang mencari event_id kembar.

Yang masih belum beres

Backpressure masih berbasis memori. Kalau konsumen melambat lebih dari beberapa menit, buffer tumbuh. Untuk volume sekarang aman, tapi ini utang yang saya catat.