06. ingest: polling
Bagian kamu: ingest/ingest.go (loop polling).
Milik Grace: ingest/source.go dan klien di internal/source/bmkg dan internal/source/pvmbg.
Syntax: belajar_go/07_goroutine_channel_select.go (goroutine, ticker, select), 05_interface.go,
06_context_defer.go.
Kenapa ada
Spek, bagian Model Ingest: BNPB mengambil data dengan polling berkala. Mock BMKG dan PVMBG pasif: mereka hanya menjawab saat ditanya, tidak mengirim sendiri. Intervalnya harus bisa diatur (disarankan 2 sampai 5 detik).
Kenapa dibutuhkan
Ini satu-satunya jalur data masuk ke sistem. Tanpa polling, tidak ada HazardEvent yang tersimpan, dipublikasikan, atau disajikan. Semua Problem (P1, P2, P4, P5) bergantung padanya.
Siapa mengerjakan apa
main.go
| menyiapkan klien dan handler
v
ingest.Run <- loop polling: ticker, cursor, goroutine per sumber (KAMU)
| memanggil
v
SourceClient.Fetch(ctx, since) <- antarmuka (GRACE)
| diimplementasikan oleh
v
klien BMKG / klien PVMBG <- ambil data + panggil mapper (GRACE)
Loop polling tidak tahu apa pun tentang JSON, kredensial, atau pemetaan. Ia hanya:
panggil Fetch, serahkan hasilnya ke handler, geser cursor.
Bagian 1: kontrak dari Grace, SourceClient
Dikutip dari internal/ingest/source.go, baris 11 sampai 14
type SourceClient interface {
Source() domain.Source
Fetch(ctx context.Context, since *time.Time) ([]domain.HazardEvent, error)
}Semua klien sumber harus punya dua method:
Source()menyebut namanya (BMKG atau PVMBG), untuk log.Fetch(ctx, since)mengambil data yang lebih baru darisincedan mengembalikanHazardEventyang sudah dipetakan.
since bertipe *time.Time (pointer): nil berarti “belum pernah mengambil, minta semuanya”.
(Pelajaran 3 dan 5.)
Dikutip dari internal/ingest/source.go, baris 24 sampai 29
type SourceError struct {
Source domain.Source
Kind ErrorKind
StatusCode int
Err error
}Kalau Fetch gagal, ia mengembalikan SourceError yang membawa jenis kegagalan:
Kind | Artinya | Contoh |
|---|---|---|
unauthorized | kredensial ditolak | 401 atau 403, kunci salah |
unavailable | sumber tidak bisa dijangkau | jaringan putus, timeout, 5xx, outage PVMBG |
invalid_response | jawaban tidak bisa dipakai | JSON rusak, record gagal dipetakan |
Perbedaan jenis ini dipakai loop untuk memutuskan tingkat log.
Bagian 2: Run, satu goroutine per sumber
Dikutip dari internal/ingest/ingest.go, baris 17 sampai 17
type Handler func(ctx context.Context, events []domain.HazardEvent) errorHandler adalah nama untuk sebuah bentuk fungsi: “fungsi yang menerima ctx dan daftar
event, dan mengembalikan error”. main.go yang mengisinya (menyimpan ke store). Dengan cara
ini ingest tidak tahu apa pun tentang store, jadi bisa diuji tanpa MongoDB. (Pelajaran 8.)
Dikutip dari internal/ingest/ingest.go, baris 21 sampai 32
func Run(ctx context.Context, logger *slog.Logger, interval time.Duration, clients []SourceClient, handle Handler) {
var wg sync.WaitGroup
for _, client := range clients {
wg.Add(1)
go func() {
defer wg.Done()
p := &poller{logger: logger, client: client, handle: handle}
p.loop(ctx, interval)
}()
}
wg.Wait()
}Untuk setiap klien, dimulai satu goroutine (pekerja yang berjalan bersamaan).
wg.Wait() menunggu semua goroutine selesai, yang terjadi saat service diminta berhenti.
Kenapa satu goroutine per sumber? Problem 2 menyebut head-of-line blocking: kalau BMKG dan
PVMBG dipanggil berurutan dalam satu loop, BMKG yang cepat ikut menunggu PVMBG yang lambat.
Dengan loop terpisah, PVMBG yang macet atau mati tidak menunda BMKG sama sekali.
Terbukti: saat outage PVMBG, event BMKG tetap bertambah (bukti_p4/06_...). (Pelajaran 7.)
Bagian 3: poller dan loop
Dikutip dari internal/ingest/ingest.go, baris 34 sampai 39
type poller struct {
logger *slog.Logger
client SourceClient
handle Handler
since *time.Time // newest OccurredAt already handled; owned by this poller's goroutine
}poller menyimpan status satu sumber. Perhatikan since: itulah cursor, yaitu “sampai
event mana saya sudah memproses”. Karena since hanya disentuh goroutine milik poller itu
sendiri, tidak perlu kunci (mutex): tidak ada dua goroutine yang menulis variabel yang sama.
Dikutip dari internal/ingest/ingest.go, baris 41 sampai 53
func (p *poller) loop(ctx context.Context, interval time.Duration) {
p.cycle(ctx)
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
p.cycle(ctx)
}
}
}Bentuknya:
p.cycle(ctx): ambil sekali langsung, tanpa menunggu satu interval dulu.- Buat ticker yang berbunyi tiap
interval(default 3 detik). - Di dalam
for,selectmenunggu salah satu:<-ctx.Done(): service diminta berhenti (docker compose stop). Loopreturn.<-ticker.C: waktunya polling lagi. Panggilcycle.
Kalau satu siklus lebih lama dari interval (misalnya PVMBG menggantung sampai timeout), ticker menjatuhkan tick yang terlewat, sehingga siklus tidak menumpuk.
Bagian 4: cycle, satu putaran
Dikutip dari internal/ingest/ingest.go, baris 56 sampai 83
func (p *poller) cycle(parent context.Context) {
ctx := core.WithCorrelationID(parent, core.NewCorrelationID())
source := string(p.client.Source())
events, err := p.client.Fetch(ctx, p.since)
if err != nil {
if ctx.Err() != nil {
return
}
p.logFetchError(ctx, source, err)
return
}
if len(events) == 0 {
p.logger.DebugContext(ctx, "poll cycle: nothing new", "source", source)
return
}
if err := p.handle(ctx, events); err != nil {
p.logger.ErrorContext(ctx, "handle events failed, will retry next cycle", "source", source, "error", err.Error())
return
}
// The cursor moves only after the handler succeeded.
if latest := newestOccurredAt(events); p.since == nil || latest.After(*p.since) {
p.since = &latest
}
p.logger.InfoContext(ctx, "poll cycle done", "source", source, "events", len(events))
}Langkah demi langkah:
- Correlation ID baru untuk siklus ini (
WithCorrelationID(parent, NewCorrelationID())). Semua log dan panggilan keluar dalam siklus ini membawa ID yang sama. Fetch(ctx, p.since): minta data yang lebih baru dari cursor.- Kalau
Fetchgagal:- Kalau penyebabnya
ctxdibatalkan (service sedang berhenti), keluar diam-diam. - Selain itu catat lewat
logFetchError, lalureturn. Cursor tidak digeser, jadi siklus berikutnya mencoba lagi dari titik yang sama.
- Kalau penyebabnya
- Kalau tidak ada event baru: catat di level debug dan selesai.
p.handle(ctx, events): serahkan ke handler (dimain.go: simpan ke Mongo). Kalau handler gagal (misalnya Mongo mati), catat ERROR danreturn. Cursor tidak digeser. Ini yang disebut retry: event yang sama diambil lagi di siklus berikutnya.- Hanya setelah handler berhasil, cursor maju ke
OccurredAtterbaru di antara event itu, dan hanya maju, tidak pernah mundur.
Kenapa cursor hanya maju setelah berhasil? Supaya tidak ada data yang hilang diam-diam. Kalau
cursor maju lebih dulu lalu penyimpanan gagal, event itu tidak akan pernah diambil lagi. Karena
store.Save adalah upsert, mengambil event yang sama dua kali aman.
Dikutip dari internal/ingest/ingest.go, baris 88 sampai 99
func (p *poller) logFetchError(ctx context.Context, source string, err error) {
attrs := []any{"source", source, "error", err.Error()}
var se *SourceError
if errors.As(err, &se) {
attrs = append(attrs, "kind", string(se.Kind))
if se.Kind == ErrorUnauthorized {
p.logger.ErrorContext(ctx, "source rejected our credentials", attrs...)
return
}
}
p.logger.WarnContext(ctx, "poll failed", attrs...)
}Tingkat log dibedakan menurut jenis kegagalan:
unauthorized→ ERROR. Kredensial salah tidak akan sembuh sendiri; perlu dilihat orang.- selain itu → WARN. Sumber mati atau rusak itu wajar, apalagi saat demo outage PVMBG.
errors.As(err, &se) memeriksa apakah error itu berjenis *SourceError dan, kalau ya, mengisi
se. (Pelajaran 2.)
Bagian 5: klien sumber (milik Grace), supaya kamu bisa menjelaskan perannya
Dikutip dari internal/source/bmkg/client.go, baris 60 sampai 92
func (client *Client) Fetch(
ctx context.Context,
since *time.Time,
) ([]domain.HazardEvent, error) {
events, err := fetch[mapping.SeismicEvent](client, ctx, "/seismic-events", since)
if err != nil {
return nil, err
}
warnings, err := fetch[mapping.TsunamiWarning](client, ctx, "/tsunami-warnings", since)
if err != nil {
return nil, err
}
warningsByEvent := make(map[string]mapping.TsunamiWarning, len(warnings))
for _, warning := range warnings {
warningsByEvent[warning.RelatedEventID] = warning
}
mapped := make([]domain.HazardEvent, 0, len(events))
for index, event := range events {
var warning *mapping.TsunamiWarning
if related, ok := warningsByEvent[event.EventID]; ok {
warningCopy := related
warning = &warningCopy
}
hazard, err := client.mapper.MapSeismicEvent(event, warning)
if err != nil {
return nil, invalidResponse(fmt.Errorf("map BMKG event at index %d: %w", index, err))
}
mapped = append(mapped, hazard)
}
return mapped, nil
}Klien BMKG:
- Ambil
/seismic-events, lalu/tsunami-warnings, dengansinceyang sama. - Indeks warning menurut
related_event_id(warningsByEvent, sebuah map). - Petakan tiap event bersama warning-nya kalau ada, lewat
MapSeismicEvent.
Kenapa digabung? Aturan severity di spek: kalau ada TsunamiWarning, warning menimpa aturan
magnitude (mis. magnitude 3 dengan warning Awas menjadi AWAS). Event dan warning datang dari
dua endpoint berbeda, jadi harus digabung sebelum dipetakan.
Klien PVMBG lebih sederhana: satu endpoint, GET /volcanic-reports, lalu MapVolcanicReport
untuk tiap laporan.
Kredensial berbeda dan tidak boleh tertukar: BMKG memakai header X-BMKG-Key, PVMBG memakai
Authorization: Bearer <token>.
Keterbatasan klien Grace (saran, bukan bug untuk data mock; sudah diperiksa terhadap data hasil polling, 9 dari 9 gempa bertsunami memiliki warning dan severity-nya cocok):
- Satu record yang gagal dipetakan menggagalkan seluruh
Fetch, sehingga satu record rusak memblokir sumber itu terus-menerus (cursor tidak maju). - Warning yang terbit sesudah event-nya sudah diproses tidak mengoreksi severity event itu. Di mock, warning selalu tersedia bersamaan dengan event-nya.
Cara menjelaskan ingest dalam 30 detik
ingest.Runmemberi satu goroutine per sumber, masing-masing mem-polling tiap interval. Setiap siklus membuat correlation ID baru, memanggilFetchpada klien sumber dengan cursorsince, lalu menyerahkan event ke handler yang menyimpannya. Cursor hanya maju setelah penyimpanan berhasil, jadi kegagalan di tengah jalan hanya berarti retry, bukan data hilang. Karena tiap sumber punya loop sendiri, PVMBG yang lambat atau mati tidak menunda BMKG.
Yang sudah dibuktikan
Terhadap bmkg-mock dan pvmbg-mock asli di docker compose (bukti_p4/):
- polling berjalan dengan correlation ID per siklus;
- outage PVMBG dicatat
WARNjenisunavailable, BMKG tetap jalan, PVMBG pulih tanpa restart dan mengejar event yang terlewat; - perubahan skema PVMBG masuk tanpa restart;
- Mongo mati lalu hidup: penyimpanan gagal diulang, tanpa duplikat.
Batasan
- Cursor hanya di memori. Setelah restart, polling mengambil semua data lagi (aman karena upsert, tetapi event lama dikirim ulang ke tahap berikutnya, jadi consumer Kafka harus idempotent).
- Mode outage
hangPVMBG belum diuji, baru modeerror(503). ingest.gotidak punya tes unit di repo (tes yang dijalankan bersifat sementara, sesuai permintaan tidak menaruh tes di repo).
Latihan
- Apa yang terjadi pada cursor kalau
handlemengembalikan error? Kenapa itu aman? - Kenapa
selectmemilikicase <-ctx.Done()? Apa yang terjadi saat kamu menjalankandocker compose stop bnpb-aggregatortanpa case itu? - Kalau kamu menambah sumber ketiga, apa saja yang harus diubah di
ingest.go? (Jawab: tidak ada. Cukup tambahkan satu klien ke daftarclientsdimain.go.)