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
- Topic — Channel/kategori untuk messages (seperti tabel, tapi append-only)
- Producer — Aplikasi yang mengirim data ke topic
- Consumer — Aplikasi yang membaca data dari topic
- Partition — Subdivisi topic untuk paralelisme
- Consumer Group — Sekumpulan consumer yang berbagi beban kerja
- Offset — Posisi baca consumer di partition (seperti bookmark)
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