Kafka for Data Pipelines — Data Engineering

Apache Kafka untuk Data Pipelines Apache Kafka adalah distributed event streaming platform yang menjadi tulang punggung data pipeline modern. Kafka berfungsi…

Apache Kafka untuk Data Pipelines

Apache Kafka adalah distributed event streaming platform yang menjadi tulang punggung data pipeline modern. Kafka berfungsi sebagai "central nervous system" yang menghubungkan semua sistem data.

Konsep Utama

Producer: Mengirim Data

from kafka import KafkaProducer
import json

producer = KafkaProducer(
    bootstrap_servers=["kafka-1:9092", "kafka-2:9092"],
    value_serializer=lambda v: json.dumps(v).encode("utf-8"),
    acks="all",            # Tunggu semua replica confirm
    retries=3,
    linger_ms=10,          # Batch messages 10ms untuk throughput
)

# Kirim event
producer.send("order-events", value={
    "order_id": 12345,
    "user_id": 678,
    "amount": 150000,
    "status": "created",
    "timestamp": "2024-06-15T10:30:00Z"
})

# Kirim dengan key (messages dengan key sama → partition sama → ordering terjamin)
producer.send("order-events",
    key=b"user-678",
    value={"order_id": 12345, "status": "paid"}
)
producer.flush()  # Pastikan semua terkirim

Consumer: Membaca Data

from kafka import KafkaConsumer
import json

consumer = KafkaConsumer(
    "order-events",
    bootstrap_servers=["kafka-1:9092"],
    group_id="warehouse-loader",       # Consumer group
    auto_offset_reset="earliest",       # Mulai dari awal jika baru
    enable_auto_commit=False,           # Manual commit untuk reliability
    value_deserializer=lambda m: json.loads(m)
)

for message in consumer:
    event = message.value
    try:
        # Proses event
        load_to_warehouse(event)
        # Commit offset setelah sukses
        consumer.commit()
    except Exception as e:
        # Jangan commit — akan retry message ini
        log_error(e, message)

Kafka di Data Pipeline

// Arsitektur tipikal
//
// Web App ──→ Kafka ──→ Spark/Flink ──→ Data Warehouse
//    │            │
//    │            ├──→ Elasticsearch (search)
//    │            ├──→ Redis (cache)
//    │            └──→ ML Service (predictions)
//    │
//    └──→ PostgreSQL (OLTP)
//              │
//              └──→ Debezium ──→ Kafka (CDC)

// Keuntungan Kafka sebagai central hub:
// 1. Decouple producer dan consumer
// 2. Buffer saat consumer lambat
// 3. Replay data (baca ulang dari offset tertentu)
// 4. Multiple consumers dari satu stream

Yang akan kamu pelajari