Tutorial 13: Apache Kafka untuk Pipeline ML Real-Time
Daftar Isi
Pendahuluan
Dalam sistem machine learning modern, kemampuan untuk memproses data dan menghasilkan prediksi secara real-time merupakan keunggulan kompetitif yang sangat penting. Pipeline inferensi batch yang berjalan setiap jam atau harian tidak mampu memenuhi kebutuhan sistem deteksi penipuan, mesin rekomendasi, penetapan harga dinamis, atau deteksi anomali yang membutuhkan waktu respons di bawah satu detik.
Apache Kafka adalah platform streaming event terdistribusi yang memungkinkan Anda membangun pipeline data real-time yang handal, skalabel, dan fault-tolerant. Ketika diintegrasikan dengan model ML, Kafka memungkinkan Anda untuk menyerap data streaming, menghitung fitur secara langsung, menyajikan prediksi real-time, dan mengirimkan hasil ke sistem hilir (downstream) — semuanya dengan throughput tinggi dan toleransi terhadap kegagalan.
Tutorial ini akan memandu Anda melalui perjalanan lengkap: dari dasar-dasar Kafka hingga membangun pipeline inferensi ML real-time tingkat produksi menggunakan Python.
Prasyarat
- Python 3.9 atau lebih tinggi
- Docker dan Docker Compose sudah terinstal
- Pemahaman dasar konsep machine learning
- Familiar dengan REST API dan JSON
- Instal paket Python yang diperlukan:
pip install confluent-kafka fastavro requests scikit-learn numpy pandas fastapi uvicorn
Memahami Apache Kafka
Konsep Inti
Apache Kafka beroperasi menggunakan beberapa abstraksi fundamental:
Topics adalah kategori atau feed bernama tempat rekaman dipublikasikan. Bayangkan topic seperti tabel database atau folder dalam sistem file. Setiap topic dibagi menjadi partitions, yaitu urutan rekaman yang teratur dan tidak dapat diubah (immutable). Partisi memungkinkan paralelisme — beberapa consumer dapat membaca dari partisi yang berbeda secara bersamaan. Producers adalah aplikasi klien yang mempublikasikan (menulis) event ke topic Kafka. Consumers adalah aplikasi yang berlangganan (membaca dan memproses) event dari topic. Consumer diorganisasikan dalam consumer groups, di mana setiap partisi dikonsumsi oleh tepat satu consumer dalam grup, memungkinkan penyeimbangan beban. Brokers adalah server Kafka yang menyimpan data dan melayani klien. Sebuah cluster Kafka terdiri dari beberapa broker untuk redundansi dan skalabilitas. Offsets adalah ID sekuensial unik yang diberikan pada setiap rekaman dalam partisi. Consumer melacak posisi mereka menggunakan offset, memungkinkan jaminan pemrosesan exactly-once atau at-least-once.Mengapa Kafka untuk ML?
| Fitur | Manfaat untuk ML |
|-------|-----------------|
| Throughput tinggi | Menangani jutaan permintaan prediksi per detik |
| Durabilitas | Tidak kehilangan data input meskipun layanan ML mati |
| Decoupling | Memisahkan penyerapan data dari inferensi model |
| Replayability | Memproses ulang data historis dengan model yang diperbarui |
| Skalabilitas | Menambahkan lebih banyak consumer/partisi seiring peningkatan beban |
Menyiapkan Kafka Secara Lokal
Buat file docker-compose.yml untuk stack Kafka lengkap:
version: '3.8'
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.5.0
environment:
ZOOKEEPERCLIENTPORT: 2181
ZOOKEEPERTICKTIME: 2000
ports:
- "2181:2181"
kafka:
image: confluentinc/cp-kafka:7.5.0
dependson:
- zookeeper
ports:
- "9092:9092"
environment:
KAFKABROKERID: 1
KAFKAZOOKEEPERCONNECT: zookeeper:2181
KAFKAADVERTISEDLISTENERS: PLAINTEXT://localhost:9092
KAFKAOFFSETSTOPICREPLICATIONFACTOR: 1
KAFKAAUTOCREATETOPICSENABLE: "true"
schema-registry:
image: confluentinc/cp-schema-registry:7.5.0
dependson:
- kafka
ports:
- "8081:8081"
environment:
SCHEMAREGISTRYHOSTNAME: schema-registry
SCHEMAREGISTRYKAFKASTOREBOOTSTRAPSERVERS: kafka:9092
kafka-connect:
image: confluentinc/cp-kafka-connect:7.5.0
dependson:
- kafka
- schema-registry
ports:
- "8083:8083"
environment:
CONNECTBOOTSTRAPSERVERS: kafka:9092
CONNECTRESTPORT: 8083
CONNECTGROUPID: "connect-cluster"
CONNECTCONFIGSTORAGETOPIC: "connect-configs"
CONNECTOFFSETSTORAGETOPIC: "connect-offsets"
CONNECTSTATUSSTORAGETOPIC: "connect-status"
CONNECTCONFIGSTORAGEREPLICATIONFACTOR: 1
CONNECTOFFSETSTORAGEREPLICATIONFACTOR: 1
CONNECTSTATUSSTORAGEREPLICATIONFACTOR: 1
CONNECTKEYCONVERTER: "org.apache.kafka.connect.json.JsonConverter"
CONNECTVALUECONVERTER: "org.apache.kafka.connect.json.JsonConverter"
Jalankan stack:
docker-compose up -d
Verifikasi semua layanan berjalan:
docker-compose ps
Producer dan Consumer dengan Python
Membangun Kafka Producer
Library confluent-kafka adalah klien Python resmi untuk Kafka. Berikut adalah producer yang mengirim data fitur ML terstruktur:
import json
import time
import random
from confluentkafka import Producer
def callbackpengiriman(err, msg):
"""Dipanggil untuk setiap pesan yang diproduksi untuk menunjukkan hasil pengiriman."""
if err is not None:
print(f"Pengiriman pesan gagal: {err}")
else:
print(f"Pesan terkirim ke {msg.topic()} [{msg.partition()}] pada offset {msg.offset()}")
def buatproducer():
"""Membuat dan mengkonfigurasi Kafka producer."""
config = {
'bootstrap.servers': 'localhost:9092',
'client.id': 'ml-feature-producer',
'acks': 'all', # Tunggu semua replika mengakui
'retries': 3, # Ulangi pada kegagalan sementara
'retry.backoff.ms': 100, # Tunggu antara pengulangan
'linger.ms': 5, # Batch pesan untuk efisiensi
'batch.size': 16384, # Ukuran batch dalam byte
'compression.type': 'snappy', # Kompresi untuk throughput
}
return Producer(config)
def buateventtransaksi():
"""Simulasi event transaksi keuangan."""
return {
'transactionid': f"txn{random.randint(100000, 999999)}",
'userid': f"user{random.randint(1, 1000)}",
'amount': round(random.uniform(1.0, 10000.0), 2),
'merchantcategory': random.choice(['grocery', 'electronics', 'restaurant', 'travel', 'atm']),
'timestamp': time.time(),
'locationlat': round(random.uniform(-90, 90), 6),
'locationlon': round(random.uniform(-180, 180), 6),
'devicetype': random.choice(['mobile', 'web', 'pos']),
}
def jalankanproducer(topic='raw-transactions', jumlahevent=100):
"""Memproduksi event transaksi simulasi ke Kafka."""
producer = buatproducer()
for i in range(jumlahevent):
event = buateventtransaksi()
key = event['userid']
value = json.dumps(event)
producer.produce(
topic=topic,
key=key,
value=value,
callback=callbackpengiriman
)
# Memicu callback pengiriman
producer.poll(0)
time.sleep(0.01)
# Tunggu semua pesan terkirim
producer.flush(timeout=30)
print(f"Memproduksi {jumlahevent} event ke topic '{topic}'")
if name == 'main':
jalankanproducer()
Membangun Kafka Consumer
import json
from confluentkafka import Consumer, KafkaError, KafkaException
def buatconsumer(groupid='ml-inference-group'):
"""Membuat dan mengkonfigurasi Kafka consumer."""
config = {
'bootstrap.servers': 'localhost:9092',
'group.id': groupid,
'auto.offset.reset': 'earliest',
'enable.auto.commit': False, # Commit manual untuk kontrol
'max.poll.interval.ms': 300000,
'session.timeout.ms': 10000,
'fetch.min.bytes': 1,
'fetch.max.wait.ms': 500,
}
return Consumer(config)
def prosespesan(message):
"""Memproses satu pesan Kafka."""
event = json.loads(message.value().decode('utf-8'))
print(f"Memproses transaksi {event['transactionid']} "
f"untuk pengguna {event['userid']}, jumlah: ${event['amount']}")
return event
def jalankanconsumer(topics=['raw-transactions']):
"""Mengkonsumsi dan memproses pesan dari Kafka."""
consumer = buatconsumer()
consumer.subscribe(topics)
try:
while True:
msg = consumer.poll(timeout=1.0)
if msg is None:
continue
if msg.error():
if msg.error().code() == KafkaError.PARTITIONEOF:
print(f"Mencapai akhir partisi {msg.partition()}")
continue
else:
raise KafkaException(msg.error())
event = prosespesan(msg)
consumer.commit(asynchronous=False)
except KeyboardInterrupt:
print("Consumer diinterupsi")
finally:
consumer.close()
if name == 'main':
jalankanconsumer()
Membangun Pipeline Fitur Streaming
Pipeline fitur streaming mengubah event mentah menjadi fitur siap ML secara real-time. Ini adalah jembatan kritis antara data mentah dan inferensi model.
import json
import time
import numpy as np
from collections import defaultdict
from confluentkafka import Consumer, Producer
class MesinFiturStreaming:
"""Rekayasa fitur real-time dari stream Kafka."""
def init(self):
self.riwayatpengguna = defaultdict(list)
self.statistikkategori = defaultdict(lambda: {'jumlah': 0, 'total': 0.0})
self.consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'feature-engine',
'auto.offset.reset': 'earliest',
'enable.auto.commit': False,
})
self.producer = Producer({
'bootstrap.servers': 'localhost:9092',
'acks': 'all',
})
def hitungfitur(self, event):
"""Mengubah event mentah menjadi fitur ML."""
userid = event['userid']
amount = event['amount']
category = event['merchantcategory']
timestamp = event['timestamp']
# Perbarui riwayat pengguna (simpan 100 transaksi terakhir)
self.riwayatpengguna[userid].append({
'amount': amount,
'timestamp': timestamp,
'category': category,
})
if len(self.riwayatpengguna[userid]) > 100:
self.riwayatpengguna[userid] = self.riwayatpengguna[userid][-100:]
riwayat = self.riwayatpengguna[userid]
jumlahlist = [h['amount'] for h in riwayat]
# Perbarui statistik kategori
self.statistikkategori[category]['jumlah'] += 1
self.statistikkategori[category]['total'] += amount
# Hitung fitur
fitur = {
'transactionid': event['transactionid'],
'userid': userid,
'amount': amount,
'amountzscore': self.hitungzscore(amount, jumlahlist),
'useravgamount': np.mean(jumlahlist),
'userstdamount': np.std(jumlahlist) if len(jumlahlist) > 1 else 0.0,
'usermaxamount': max(jumlahlist),
'userminamount': min(jumlahlist),
'usertxncount': len(jumlahlist),
'timesincelasttxn': self.waktusejakterakhir(riwayat, timestamp),
'categoryavgamount': (
self.statistikkategori[category]['total'] /
self.statistikkategori[category]['jumlah']
),
'amounttocategoryratio': (
amount / (self.statistikkategori[category]['total'] /
self.statistikkategori[category]['jumlah'])
if self.statistikkategori[category]['jumlah'] > 0 else 1.0
),
'isnewuser': 1 if len(jumlahlist) <= 2 else 0,
'hourofday': int((timestamp % 86400) / 3600),
'devicetypeencoded': {'mobile': 0, 'web': 1, 'pos': 2}.get(
event.get('devicetype', 'web'), 1
),
}
return fitur
def hitungzscore(self, nilai, riwayat):
if len(riwayat) < 2:
return 0.0
mean = np.mean(riwayat)
std = np.std(riwayat)
return (nilai - mean) / std if std > 0 else 0.0
def waktusejakterakhir(self, riwayat, waktusaatini):
if len(riwayat) < 2:
return 0.0
return waktusaatini - riwayat[-2]['timestamp']
def jalankan(self, inputtopic='raw-transactions', outputtopic='ml-features'):
"""Loop utama: konsumsi event mentah, hitung fitur, produksi ke output."""
self.consumer.subscribe([inputtopic])
print(f"Mesin fitur dimulai: {inputtopic} -> {outputtopic}")
try:
while True:
msg = self.consumer.poll(timeout=1.0)
if msg is None or msg.error():
continue
event = json.loads(msg.value().decode('utf-8'))
fitur = self.hitungfitur(event)
self.producer.produce(
topic=outputtopic,
key=event['userid'],
value=json.dumps(fitur),
)
self.producer.poll(0)
self.consumer.commit(asynchronous=False)
except KeyboardInterrupt:
pass
finally:
self.consumer.close()
self.producer.flush()
if name == 'main':
mesin = MesinFiturStreaming()
mesin.jalankan()
Inferensi ML Real-Time dengan Kafka
Sekarang kita membangun layanan inferensi yang mengkonsumsi vektor fitur dan menghasilkan prediksi:
import json
import pickle
import numpy as np
from confluentkafka import Consumer, Producer
class InferensiMLRealTime:
"""Mengkonsumsi fitur dari Kafka, menjalankan inferensi ML, memproduksi prediksi."""
KOLOMFITUR = [
'amount', 'amountzscore', 'useravgamount', 'userstdamount',
'usermaxamount', 'userminamount', 'usertxncount',
'timesincelasttxn', 'categoryavgamount',
'amounttocategoryratio', 'isnewuser', 'hourofday',
'devicetypeencoded',
]
def init(self, modelpath='fraudmodel.pkl'):
with open(modelpath, 'rb') as f:
self.model = pickle.load(f)
print(f"Model dimuat dari {modelpath}")
self.consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'ml-inference-service',
'auto.offset.reset': 'latest',
'enable.auto.commit': False,
})
self.producer = Producer({
'bootstrap.servers': 'localhost:9092',
'acks': 'all',
})
def prediksi(self, fitur):
"""Menjalankan inferensi model pada vektor fitur."""
vektorfitur = np.array(
[fitur.get(col, 0.0) for col in self.KOLOMFITUR]
).reshape(1, -1)
prediksi = self.model.predict(vektorfitur)[0]
probabilitas = self.model.predictproba(vektorfitur)[0]
return {
'transactionid': fitur['transactionid'],
'userid': fitur['userid'],
'isfraud': int(prediksi),
'fraudprobability': float(max(probabilitas)),
'modelversion': 'v1.0.0',
}
def jalankan(self, inputtopic='ml-features', outputtopic='predictions'):
"""Loop inferensi utama."""
self.consumer.subscribe([inputtopic])
print(f"Layanan inferensi dimulai: {inputtopic} -> {outputtopic}")
try:
while True:
msg = self.consumer.poll(timeout=1.0)
if msg is None or msg.error():
continue
fitur = json.loads(msg.value().decode('utf-8'))
hasil = self.prediksi(fitur)
self.producer.produce(
topic=outputtopic,
key=fitur['userid'],
value=json.dumps(hasil),
)
self.producer.poll(0)
self.consumer.commit(asynchronous=False)
if hasil['isfraud']:
print(f"PERINGATAN PENIPUAN: {hasil['transactionid']} "
f"(prob: {hasil['fraudprobability']:.4f})")
except KeyboardInterrupt:
pass
finally:
self.consumer.close()
self.producer.flush()
Schema Registry untuk Kontrak Data
Schema Registry menegakkan kontrak pada data yang mengalir melalui topic Kafka, mencegah ketidakcocokan skema yang dapat merusak pipeline ML Anda.
import json
from confluentkafka import Producer
from confluentkafka.serialization import (
SerializationContext, MessageField
)
from confluentkafka.schemaregistry import SchemaRegistryClient
from confluentkafka.schemaregistry.avro import AvroSerializer
Definisi skema Avro untuk event transaksi
SKEMATRANSAKSI = json.dumps({
"type": "record",
"name": "Transaction",
"namespace": "com.ml.fraud",
"fields": [
{"name": "transactionid", "type": "string"},
{"name": "userid", "type": "string"},
{"name": "amount", "type": "double"},
{"name": "merchantcategory", "type": "string"},
{"name": "timestamp", "type": "double"},
{"name": "locationlat", "type": "double"},
{"name": "locationlon", "type": "double"},
{"name": "devicetype", "type": "string"},
]
})
def buatavroproducer():
"""Membuat producer dengan serialisasi Avro dan Schema Registry."""
klienregistry = SchemaRegistryClient({
'url': 'http://localhost:8081'
})
serializeravro = AvroSerializer(
klienregistry,
SKEMATRANSAKSI,
lambda obj, ctx: obj,
)
producer = Producer({
'bootstrap.servers': 'localhost:9092',
})
return producer, serializeravro
Evolusi Skema
Schema Registry mendukung evolusi skema dengan mode kompatibilitas:
- BACKWARD: Skema baru dapat membaca data yang diproduksi dengan skema sebelumnya.
- FORWARD: Skema sebelumnya dapat membaca data yang diproduksi dengan skema baru.
- FULL: Kompatibel baik ke belakang maupun ke depan.
# Mengatur tingkat kompatibilitas untuk sebuah subject
curl -X PUT -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data '{"compatibility": "BACKWARD"}' \
http://localhost:8081/config/raw-transactions-value
Kafka Connect untuk Integrasi Data
Kafka Connect memindahkan data antara Kafka dan sistem eksternal tanpa menulis kode. Ini berguna untuk menyimpan prediksi ke database atau mengambil data pelatihan.
Menyimpan Prediksi ke PostgreSQL
{
"name": "predictions-postgres-sink",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
"tasks.max": "2",
"topics": "predictions",
"connection.url": "jdbc:postgresql://postgres:5432/mlresults",
"connection.user": "mluser",
"connection.password": "mlpassword",
"auto.create": "true",
"auto.evolve": "true",
"insert.mode": "upsert",
"pk.mode": "recordvalue",
"pk.fields": "transactionid"
}
}
Deploy konektor:
curl -X POST -H "Content-Type: application/json" \
--data @predictions-sink.json \
http://localhost:8083/connectors
Praktik Terbaik untuk Produksi
1. Strategi Partisi
Pilih kunci partisi yang mendistribusikan beban secara merata sambil mempertahankan jaminan urutan jika diperlukan. Untuk inferensi ML, partisi berdasarkan userid memastikan semua event untuk satu pengguna diproses secara berurutan.
2. Penskalaan Consumer Group
Skalakan consumer secara horizontal dengan menambahkan lebih banyak instance ke consumer group yang sama. Setiap partisi ditugaskan ke tepat satu consumer.
3. Penanganan Error dan Dead Letter Queue
Jangan pernah kehilangan pesan — arahkan pemrosesan yang gagal ke dead letter queue:
def prosesdengandlq(consumer, producer, msg):
try:
event = json.loads(msg.value().decode('utf-8'))
hasil = proses
event(event)
producer.produce(topic='predictions', value=json.dumps(hasil))
except Exception as e:
rekamanerror = {
'pesanasli': msg.value().decode('utf-8'),
'error': str(e),
'topic': msg.topic(),
'partisi': msg.partition(),
'offset': msg.offset(),
}
producer.produce(
topic='ml-inference-dlq',
value=json.dumps(rekamanerror),
)
4. Monitoring dan Observabilitas
Lacak metrik utama untuk pipeline ML Kafka Anda:
- Consumer lag: Seberapa tertinggal consumer dari pesan terbaru.
- Throughput: Pesan yang diproses per detik.
- Latensi: Waktu end-to-end dari produksi event hingga output prediksi.
- Tingkat error: Persentase pesan yang diarahkan ke DLQ.
5. Pemrosesan Idempoten
Pastikan layanan inferensi Anda dapat memproses ulang pesan yang sama tanpa efek samping:
import redis
klienredis = redis.Redis(host='localhost', port=6379, db=0)
def prosesidempoten(event):
txnid = event['transactionid']
if klienredis.exists(f"diproses:{txnid}"):
return None # Sudah diproses, lewati
hasil = jalankaninferensi(event)
klienredis.setex(f"diproses:{txnid}", 3600, "1") # TTL 1 jam
return hasil
Kesimpulan
Apache Kafka menyediakan tulang punggung untuk sistem ML real-time dengan memisahkan produsen data dari konsumen, memastikan durabilitas, dan memungkinkan skalabilitas horizontal. Dalam tutorial ini, Anda telah mempelajari cara:
- Menyiapkan stack Kafka lengkap dengan Schema Registry dan Kafka Connect
- Membangun producer dan consumer Python menggunakan library
confluent-kafka - Mengimplementasikan pipeline rekayasa fitur streaming
- Men-deploy layanan inferensi ML real-time
- Menegakkan kontrak data dengan skema Avro
- Menerapkan praktik terbaik produksi termasuk penanganan error, monitoring, dan pemrosesan idempoten
Arsitektur yang disajikan di sini — event mentah mengalir melalui rekayasa fitur ke inferensi model dengan prediksi yang dikirim ke sistem hilir — adalah pola standar yang digunakan oleh perusahaan seperti Netflix, Uber, dan LinkedIn untuk platform ML real-time mereka. Mulailah dengan fondasi ini dan kembangkan sesuai kebutuhan spesifik Anda.