Streaming (Kafka) Advanced¶
🔄 Data Engineering · Level 5
When you'd use this
Real-time data pipelines with Kafka, producers, consumers and stream processing.
Handle continuous event streams (Kafka) in real time — for live analytics, pipelines, and event-driven systems.
Kafka producer¶
Publish events to a topic for other systems to consume.
from kafka import KafkaProducer
import json
producer = KafkaProducer(
bootstrap_servers=["localhost:9092"],
value_serializer=lambda v: json.dumps(v).encode("utf-8"),
)
# Send events
for i in range(100):
event = {"user_id": i, "action": "page_view", "timestamp": "2026-08-23T12:00:00"}
producer.send("user-events", value=event)
print(f" Sent event {i}")
producer.flush() # ensure all messages are sent
producer.close()
Kafka consumer¶
Read events from a topic, tracking offsets for reliable processing.
from kafka import KafkaConsumer
import json
consumer = KafkaConsumer(
"user-events",
bootstrap_servers=["localhost:9092"],
group_id="my-consumer-group",
auto_offset_reset="earliest",
value_deserializer=lambda m: json.loads(m.decode("utf-8")),
)
print("Listening for events...")
for message in consumer:
event = message.value
print(f" Received: {event['action']} from user {event['user_id']}")
# Process event (store, transform, trigger actions)
Stream processing pattern¶
Transform and aggregate continuous event streams in real time.
import asyncio
from kafka import KafkaConsumer, KafkaProducer
import json
def process_stream():
"""Read from one topic, process, write to another."""
consumer = KafkaConsumer("raw-events", bootstrap_servers=["localhost:9092"],
value_deserializer=lambda m: json.loads(m.decode()))
producer = KafkaProducer(bootstrap_servers=["localhost:9092"],
value_serializer=lambda v: json.dumps(v).encode())
for message in consumer:
event = message.value
# Transform
enriched = {
**event,
"processed_at": datetime.utcnow().isoformat(),
"is_premium": event.get("total_spend", 0) > 1000,
}
# Write to enriched topic
producer.send("enriched-events", value=enriched)
Practice Exercises¶
- Build a producer that generates fake events at 100/second.
- Build a consumer that aggregates events into 1-minute windows.
- Implement exactly-once semantics using transactions.
- Build a dead-letter queue — failed messages go to a separate topic.
- Monitor lag — track how far behind the consumer is.
💬 Discussion
Have a question about this topic? Found an error? Share your thoughts below.