๐๏ธ Advanced Streaming Architectures
๐ช Cross-Cluster Mirroring (MirrorMaker)
What Existed Previously: A company builds a single massive Kafka cluster in their US-East data center. All global traffic routes there.
Problems Faced: If a hurricane takes out the US-East data center, the entire global streaming pipeline halts. Data is lost, and applications fail.
How Present Technology Solves It: Using Kafka MirrorMaker, you can continuously replicate (mirror) data from your primary US-East cluster to a standby US-West cluster.
- High Availability (HA): If the primary cluster dies, clients automatically failover to the backup cluster.
- RPO (Recovery Point Objective): MirrorMaker minimizes RPO by ensuring the backup cluster is always just milliseconds behind the primary.
๐ฉ๏ธ Edge Cluster Aggregation
Real-World Analogy Mapping: Imagine a global retail chain like Walmart.
- Sending every single receipt instantly to the central cloud HQ is extremely expensive and causes massive network congestion.
- Instead, they put a small "Edge Cluster" (a mini Kafka cluster) inside every physical store.
- The Edge Cluster processes and aggregates the local store data, and only sends summary reports (e.g., total hourly sales) to the Central Cloud Cluster.
Edge Computing moves processing as close to the data source as possible to reduce latency and bandwidth.
๐ก IoT Use Case: Faulty Cell Tower Detection
Let's look at an end-to-end predictive maintenance pipeline.
Follow the Data:
- IoT Sensors: Thousands of cell towers emit real-time telemetry (temperature, voltage, signal strength) every second.
- Kafka Ingestion: The telemetry is published to a high-throughput Kafka topic called
tower-metrics. - Spark Structured Streaming: Spark reads directly from the Kafka topic using the Direct Stream approach.
- Machine Learning Transformation: Spark applies a pre-trained ML model to the incoming metrics to detect anomalies (e.g., voltage spikes that precede a hardware failure).
- Alerting Sink: If the ML model flags an anomaly, Spark immediately writes the alert to a Kafka
maintenance-alertstopic. - Action: An automated system reads the alert and dispatches a repair crew before the tower actually breaks down.
By combining Kafka's massive ingestion capabilities with Spark's real-time ML processing, telecom companies eliminate reactive user complaints and replace them with proactive predictive maintenance.
๐ป Practical: Filtering Kafka ERROR Logs
The Smart Way: Here is a Structured Streaming job that continuously consumes application logs directly from Kafka, filters for "ERROR", and counts them using a 1-minute tumbling window.
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, window
spark = SparkSession.builder.appName("KafkaErrorWindow").getOrCreate()
# 1. Read directly from Kafka
source = spark.readStream.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "logs") \
.load()
# 2. Extract raw bytes into a string
lines = source.selectExpr("CAST(value AS STRING) as line", "timestamp")
# 3. Filter for ERROR logs only
errors = lines.filter(col("line").contains("ERROR"))
# 4. Group by a 1-minute window and count
agg = errors.groupBy(window(col("timestamp"), "1 minute")).count()
# 5. Output results directly to the console
query = agg.writeStream.outputMode("complete").format("console").start()
query.awaitTermination()
๐งช Practice Drill
// Try answering these:
Q1. What Kafka tool is used to replicate data between two entirely different data centers for Disaster Recovery?
Q2. Why use "Edge Clusters" instead of just sending all raw data directly to a central cloud cluster?
Q3. In the cell tower use case, what is the core benefit of processing the data in real-time streams instead of nightly batches?
๐ก Click for Solutions
A1. MirrorMaker.
A2. To reduce network bandwidth and latency. Edge clusters process the heavy raw data locally and only send lightweight, aggregated summaries to the central cloud.
A3. Predictive Maintenance. Real-time processing allows you to detect the anomaly and dispatch a technician before the hardware actually fails, avoiding downtime and customer complaints.
โ ๐ฅ Kafka Consumer API & Monitoring | Syllabus โ ๐ Kafka & Spark Streaming Syllabus