Real-time Analytics Pipeline
Untuk use case yang membutuhkan analitik dalam hitungan detik (live dashboard, alerting, fraud detection), kamu butuh database dan arsitektur yang didesain khusus untuk real-time ingestion dan query.
Database untuk Real-time Analytics
| Database | Kelebihan | Use Case |
|---|---|---|
| ClickHouse | Query sangat cepat, open-source | Analytics dashboard, log analysis |
| Apache Druid | Sub-second OLAP, real-time ingestion | User-facing analytics |
| TimescaleDB | Time-series di atas PostgreSQL | Metrics, IoT, monitoring |
| Apache Pinot | LinkedIn-scale, low latency | User-facing real-time dashboard |
ClickHouse
Columnar database yang bisa meng-query miliaran baris dalam milidetik.
-- Create table dengan MergeTree engine
CREATE TABLE events (
event_id UUID,
user_id UInt64,
event_type LowCardinality(String),
page_url String,
created_at DateTime
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(created_at)
ORDER BY (event_type, created_at);
-- Insert data (ClickHouse optimal untuk batch insert)
INSERT INTO events VALUES
(generateUUIDv4(), 123, 'page_view', '/home', now()),
(generateUUIDv4(), 456, 'click', '/products', now());
-- Query: real-time dashboard data
SELECT
toStartOfMinute(created_at) AS minute,
event_type,
count() AS event_count,
uniqExact(user_id) AS unique_users
FROM events
WHERE created_at >= now() - INTERVAL 1 HOUR
GROUP BY minute, event_type
ORDER BY minute DESC;
Materialized Views
Pre-compute aggregasi saat data masuk, bukan saat query. Sangat powerful untuk real-time dashboard.
-- ClickHouse materialized view
-- Otomatis aggregate saat INSERT, bukan saat SELECT
CREATE MATERIALIZED VIEW hourly_stats
ENGINE = SummingMergeTree()
ORDER BY (hour, event_type)
AS SELECT
toStartOfHour(created_at) AS hour,
event_type,
count() AS event_count,
uniq(user_id) AS unique_users
FROM events
GROUP BY hour, event_type;
-- Query materialized view → instant response
SELECT * FROM hourly_stats
WHERE hour >= now() - INTERVAL 24 HOUR;
-- PostgreSQL materialized view (batch refresh)
CREATE MATERIALIZED VIEW mv_daily_revenue AS
SELECT DATE(created_at) AS date, SUM(amount) AS revenue
FROM orders GROUP BY DATE(created_at);
-- Refresh (harus manual/scheduled)
REFRESH MATERIALIZED VIEW CONCURRENTLY mv_daily_revenue;
Real-time Pipeline Architecture
// End-to-end real-time analytics
//
// App ──→ Kafka ──→ Flink/Spark Streaming ──→ ClickHouse ──→ Dashboard
// │
// (enrich, filter,
// window aggregate)
//
// Latency: event terjadi → muncul di dashboard: 1-10 detik
//
// Simpler alternative (tanpa stream processor):
// App ──→ Kafka ──→ ClickHouse (Kafka engine) ──→ Dashboard
// Latency: 1-5 detik, tapi tanpa complex transformasi