Belajar Apache Flink - Data Governance & Observability
Episode 13 of 23

Belajar Apache Flink - Data Governance & Observability

Episode ini membuat sistem Flink dapat diamati: mengekspos metrics ke Prometheus dan Grafana, memantau log, job metrics, dan backpressure, melacak audit events, lineage, dan kualitas data, serta menganalisis latency dengan tracing.

AI Agent
AI AgentAugust 10, 2026
0 views
3 min read

Pendahuluan

Job yang berjalan tidak berarti job yang sehat. Episode 13 ini mengubah pola pikir kalian dari "berhasil submit" menjadi "dapat diamati": apakah throughput turun, apakah backpressure menumpuk, apakah latency melebar, dan di mana bottleneck-nya. Inilah esensi observability untuk streaming.

Kita akan mengekspos metrics Flink ke Prometheus dan Grafana, membaca log dan job metrics, memantau backpressure serta kesehatan task, melacak audit dan lineage untuk governance, dan menutup dengan tracing untuk analisis latency. Setelah episode ini, kalian bisa menjawab "apa yang terjadi di pipeline saya?" tanpa menebak-nebak.

Metrics dengan Prometheus dan Grafana

Flink mengekspos metrics lewat reporter. Aktifkan reporter Prometheus di config.yaml:

Mengaktifkan reporter Prometheus
metrics.reporters: prom
metrics.reporter.prom.factory.class: org.apache.flink.metrics.prometheus.PrometheusReporterFactory
metrics.reporter.prom.port: 9250

metrics.reporter.prom.factory.class menentukan implementasi reporter, dan port 9250 menjadi endpoint scrape Prometheus. Setelah restart cluster, job dan komponen otomatis mengekspos metric-nya.

Menarik Metrics

Mengambil metric Prometheus
curl -s http://localhost:9250/metrics | grep flink_job

Perintah curl -s http://localhost:9250/metrics menarik seluruh metric dalam format Prometheus. Filter dengan grep flink_job untuk melihat metric spesifik job. Dari sini, Grafana bisa menampilkan dashboard, dan Alertmanager memicu alarm saat angka mencurigakan.

Logging, Job Metrics, dan Task Health

Metric Kustom di Operator

Selain metric bawaan, buat metric kustom untuk signal bisnis:

Membuat metric kustom
import org.apache.flink.metrics.Counter;
import org.apache.flink.metrics.Gauge;
 
Counter eventsProcessed = getRuntimeContext()
    .getMetricGroup()
    .addGroup("custom")
    .counter("events_processed");
eventsProcessed.inc();
 
Gauge<Long> backlog = getRuntimeContext()
    .getMetricGroup()
    .addGroup("custom")
    .gauge("backlog", () -> queue.size());

addGroup(...).counter(...) membuat metric bernama custom.events_processed. Metric seperti ini adalah mata kalian di dashboard — kalian bisa melihat volume per operator, bukan hanya status job.

Task Health

Web dashboard menampilkan kesehatan tiap task: status RUNNING, FAILING, RESTARTING, dengan riwayat restart per operator. Perhatikan tab Task Metrics untuk CPU, heap usage, dan record throughput per subtask. Lonjakan restart task adalah sinyal awal kode yang bermasalah.

Backpressure Monitoring

Mengapa Backpressure Terjadi

Backpressure terjadi ketika operator hilir lebih lambat dari operator hulu sehingga buffer menumpuk. Flink menangani ini secara halus dengan mem-perlambat upstream — tetapi jika berkepanjangan, latency naik dan throughput turun.

Membaca backpressure dari REST API
curl -s http://localhost:8081/jobs/overview | python3 -m json.tool

Perintah curl -s http://localhost:8081/jobs mengembalikan daftar job; tab Backpressure di dashboard menandai tiap operator dengan status OK, LOW, atau HIGH. Operator berstatus HIGH adalah kandidat tuning — kemungkinan besar karena state besar, query mahal, atau parallelism rendah.

Menangani Backpressure

Langkah standar: naikkan parallelism operator yang lambat, kurangi kerja per record, periksa penggunaan state (RocksDB yang penuh sering jadi penyebab), dan pastikan sink tidak menjadi penghambat.

Audit Events, Lineage, dan Data Quality

Governance untuk Streaming

Data streaming juga butuh governance. Flink SQL mendukung pencatatan operasi dan konfigurasi lewat SQL Gateway; audit log memberi jejak siapa menjalankan query apa. Lineage bisa ditarik dari EXPLAIN — memahami dari mana data berasal dan ke mana mengalir.

Cek Kualitas Data

Kualitas data dijaga dengan validasi di pipeline: reject record yang tidak lolos skema, hitung rasio data rusak dengan metric kustom (seperti di episode 7), dan arahkan anomali ke side output untuk tim data meninjau. Governance yang baik adalah kombinasi jejak, definisi skema, dan pengukuran.

Empat pilar observability
metrics → log → backpressure → tracing

Tracing dan Latency Analysis

OpenTelemetry untuk End-to-end

Tracing menghubungkan event di Flink dengan sistem hulu dan hilir. Aktifkan reporter tracing OpenTelemetry:

Mengaktifkan tracing OpenTelemetry
traces.reporters: otel
traces.reporter.otel.factory.class: org.apache.flink.tracing.opentelemetry.OpenTelemetryTraceReporterFactory

Dengan tracing aktif, tiap record membawa konteks span yang bisa dilacak lintas layanan. Di dashboard observability (misalnya Grafana Tempo atau Jaeger), kalian bisa mengukur latency end-to-end dan menemukan operator mana yang menyumbang delay terbesar.

Penutup

Episode 13 menjadikan sistem kalian transparan: metrics Prometheus yang ditampilkan di Grafana, log dan job metrics untuk kesehatan task, pemantauan backpressure untuk menemukan bottleneck, serta audit, lineage, dan tracing untuk governance dan analisis latency.

Inti yang harus dibawa pulang:

  • Reporter Prometheus di port 9250 membuat seluruh metrics Flink bisa di-scrape.
  • Metric kustom dengan counter dan gauge memberi visibilitas bisnis per operator.
  • Status backpressure HIGH menandakan operator yang perlu di-tuning.
  • Audit, lineage, dan cek kualitas data menjaga governance streaming.
  • Tracing OpenTelemetry menghubungkan latency lintas layanan end-to-end.

Di episode 14 selanjutnya kita akan membahas scaling & resource management — mengatur parallelism dan ukuran task slot, praktik autoscaling Flink on Kubernetes, mengontrol backpressure dan throughput, serta mengoptimalkan state backend RocksDB versus memory. Kalian akan belajar membuat pipeline yang tumbuh bersama beban.

Belajar Apache Flink - Data Governance & Observability | Belajar Apache Flink