Real-time Analytics Pipeline — Data Engineering

Real-time Analytics Pipeline Untuk use case yang membutuhkan analitik dalam hitungan detik (live dashboard, alerting, fraud detection), kamu butuh database dan

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

DatabaseKelebihanUse Case
ClickHouseQuery sangat cepat, open-sourceAnalytics dashboard, log analysis
Apache DruidSub-second OLAP, real-time ingestionUser-facing analytics
TimescaleDBTime-series di atas PostgreSQLMetrics, IoT, monitoring
Apache PinotLinkedIn-scale, low latencyUser-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

Yang akan kamu pelajari