🧠 Chapter 11 — Akka & Kafka Cheat Sheet
Big Picture: Data flows from sources → gets ingested → gets processed (batch or real-time) → stored for querying.
Stack: Kafka (ingest) + Hadoop (batch) + Spark (processing) + NoSQL/SQL (service)
📐 Part 1 — Three Models of Data Flow
Before diving into actors, understand how data moves between services:
| Model | Unit | Response | Blocking? | Sync? | Addresses Recipient? |
|---|---|---|---|---|---|
| Databases | Data | None | Non-blocking | Async | ❌ No |
| Services (RPC) | Function call | Yes | Blocking | Sync | ✅ Yes |
| Message-Passing | Messages | Maybe | Non-blocking | Async | ✅ Yes |
🍵 Analogy
| Model | Real Life |
|---|---|
| Database | Whiteboard in the office — anyone reads/writes, no direct recipient |
| Service (RPC) | Phone call — you wait until the other person picks up and responds |
| Message-Passing | LINE message — fire-and-forget, they reply when they're ready |
Key Terms Demystified
- Asynchronous = ส่งแล้วไม่ต้องรอ → ทำงานอื่นต่อได้ (เหมือนส่ง LINE แล้วไปกินข้าวก่อน)
- Blocking = Thread ถูกหยุดจนกว่าจะได้ผลลัพธ์ → เหมือนยืนต่อแถวซื้อกาแฟ ทำอะไรไม่ได้
- Non-blocking = ยิงคำขอแล้วทำงานอื่นต่อ แล้วค่อยรับผลทีหลัง → เหมือนสั่งกาแฟผ่านแอปแล้วไปนั่งรอ
- Addressing Recipient = ระบุปลายทางชัดเจน (ส่งพัสดุระบุชื่อ-ที่อยู่) ≠ ทิ้งไว้บนโต๊ะกลางให้ใครมาหยิบ
🎭 Part 2 — Actor Model
Why Not Just Use Threads?
Writing concurrent programs with threads + locks is:
- Error-prone (race conditions, deadlocks)
- Leads to spaghetti code
- Hard to test and maintain
Solution → Actor Model (originally from Carl Hewitt's 1973 paper; popularized by Erlang at Ericsson)
Core Concept
Actor = a computational entity that treats message-passing as the universal primitive of concurrent computation.
Each actor:
- Has private state (no shared state between actors → no race conditions!)
- Owns exactly one mailbox (message queue)
- Processes one message at a time
📮 Analogy: Actors are like P.O. Boxes. You drop a letter in, the person checks it when they're ready, and replies on their own time. Nobody stands at the door waiting.
What Can an Actor Do When It Receives a Message?
- Send messages to other actors (non-blocking)
- Create new child actors (spawn!)
- Modify its own internal state for the next message
How Communication Works (Step-by-Step)
1. Sender adds message → end of recipient's queue (mailbox)
2. If not already scheduled → mark actor as "ready to execute"
3. Scheduler picks the actor → starts executing
4. Actor picks message from front of queue
5. Actor updates internal state + sends messages to others
6. Actor is unscheduled (done for this message)
Key Properties
- Messages always asynchronous (sender never waits)
- No guaranteed ordering of message arrival
- Enqueue/dequeue operations are atomic → no race conditions!
When to Use Actor Model?
- Problem can be split into independent tasks
- Problem is a pipeline of linked tasks
- In short: when the problem is parallelizable
- Remember Amdahl's Law: speedup is bounded by the sequential portion
Actor Model vs. OOP vs. Task Programming
| Object-Oriented | Actor | Task-Oriented | |
|---|---|---|---|
| Core Unit | Object | Actor | Task |
| Communication | Method calls | Async messages | Scheduled tasks |
| Key Benefit | Separation of concerns | OO + async | Async/parallel decoupling |
Advantages vs. Drawbacks
| ✅ Advantages | ❌ Drawbacks |
|---|---|
| No explicit locks → no deadlocks | Creating many actors → performance/latency cost |
| Fault-tolerant "let it crash" philosophy | Need to track where actors live (location overhead) |
| Dynamic parallelism via spawning | Active actors = computation + storage cost |
| Extends OOP benefits |
Advantages Over Pure RPC
- Fault-tolerance: "Let it crash!" — actors restart automatically; failed messages are re-routed
- Deadlock prevention: Async messaging + private state eliminates most parallelism bugs
- Parallelism: Multiple actors = multiple workers processing concurrently; actors can spawn more if needed
🚀 Part 3 — Akka (Actor Framework for JVM)
Akka is an open-source framework for Java/Scala on the JVM. Written in Scala, inspired by Erlang, created by Jonas Bonér, maintained by Lightbend.
Key Features
- Simple, high-level abstractions for concurrency
- Event-driven, asynchronous, non-blocking
- Supports actor hierarchies → great for fault-tolerant distributed apps
- Deploy actors across multiple JVMs
Akka Modules
| Module | Purpose |
|---|---|
| Akka Actors | Core: concurrency + distribution |
| Akka Cluster | Resilient distribution across multiple nodes |
| Akka Streams | Async, non-blocking, backpressured reactive streams |
| Akka Http | Streaming-first HTTP server/client |
| Cluster Sharding | Reference actors by identity, not location |
| Akka Persistence | Persist actor state for fault tolerance (≠ ephemeral!) |
| Distributed Data | Eventually consistent key-value store |
| Alpakka | Connectors to other tech stacks |
Akka Actor Internals
- Messages are immutable, serializable objects (never use mutable messages!)
- Receiver interprets a message via pattern matching
- Communication: send to mailbox → unblocking, fire-and-forget
Java Code Example
public class Worker extends AbstractActor {
@Override
public Receive createReceive() {
return receiveBuilder()
.match(String.class, this::respondTo) // match String messages
.matchAny(o -> System.out.println("Unknown!")) // fallback
.build();
}
private void respondTo(String message) {
System.out.println(message);
// Reply to sender — async, non-blocking
this.sender().tell("Got your message, thanks!", this.self());
}
}| Code Part | Meaning |
|---|---|
AbstractActor | Inherit actor behavior, state, mailbox |
Receive | Pattern matching + deserialization |
.build() | Builder pattern for constructing Receive |
this.sender().tell(...) | Send async response back to sender |
🏗️ Part 4 — Actor Hierarchies & Lifecycle
Hierarchy Structure
Actors form a tree:
Root Actor (Guardian)
└── Parent Actor
├── Child Actor A
└── Child Actor B
└── Grandchild Actor
- Creating actor = parent → supervises children
- If child fails, parent decides what to do
Fault-Tolerance Options (When Child Fails)
| Action | What Happens |
|---|---|
| restart | Kill and restart the child actor (fresh state) |
| resume | Ignore the failure, continue processing |
| stop | Permanently stop the child actor |
| escalate | Pass failure up to the parent to decide |
⚠️ If a parent fails → system crash (no standby parent!)
💡 This is the key to fault-tolerance: hierarchical actor networks let failed work be reassigned to surviving actors automatically.
Actor Lifecycle
actorOf() ──→ [Started]
│
┌───────┼───────────┐
restart() resume() stop()
│ │ │
[Restarting] (keep) [Stopped] ──→ DeadLetterBox
(catches messages to dead actors)
- ActorRef stays valid even after the actor stops (grey = dead reference, messages → DeadLetterBox)
📬 Part 5 — Message Delivery Guarantees
| Guarantee | Delivery | Duplicates? | Performance | How |
|---|---|---|---|---|
| at-most-once | 0 or 1 times | Never | 🟢 Fastest | Fire-and-forget |
| at-least-once | 1+ times | Possible | 🟡 OK | Send + acknowledge |
| exactly-once | Exactly 1 time | Never | 🔴 Slowest | Send + ack + deduplicate |
Analogy:
- at-most-once = letter with no tracking
- at-least-once = registered mail (resends if unconfirmed)
- exactly-once = notarized delivery with duplicate check
With TCP, Akka nearly guarantees exactly-once, but failures can still cause message loss.
⬆️⬇️ Part 6 — Push vs. Pull Propagation
| Model | How | Risk |
|---|---|---|
| Push | Producer immediately sends work to consumers (queued in consumer's inbox) | ⚠️ Congestion / message drops if consumer is slow |
| Pull | Consumer asks for work when ready (queued in producer's state) | ✅ No congestion, but slower propagation |
Real-world: Kafka uses a pull model — consumers poll the broker at their own pace.
🐘 Part 7 — Apache Kafka
Kafka is a distributed streaming platform — the de facto standard for near real-time data ingestion.
Why Kafka?
- High Scalability via partitioning
- Fault Tolerant via replication
- High Parallelism and decoupling between producers and consumers
Who Uses It?
| Company | Usage |
|---|---|
| 400 nodes, 18k topics, 220B msg/day (peak 3.2M msg/s) | |
| Netflix | Real-time monitoring and event processing |
| Part of Storm real-time pipelines | |
| Spotify | Log delivery: reduced from 4h → 10 seconds |
| Uber | GPS logs, surge pricing, driver/rider events |
🧱 Part 8 — Kafka Core Concepts
| Concept | Description |
|---|---|
| Broker | A single Kafka node (server) in the cluster |
| Topic | A named stream/category of records (like a channel) |
| Producer | Pushes messages into a topic |
| Consumer | Pulls messages from a topic |
| Data Retention | Messages kept based on time or size limit |
| Zookeeper | Stores cluster metadata + consumer offsets |
🎬 YouTube Analogy:
- Topic = YouTube channel
- Producer = Content creator uploading videos
- Consumer = Subscriber watching
- Broker = YouTube server
- Partitions = separate video queues per category
Kafka Broker
- Handles all read/write requests from clients
- Acts as intermediary:
Producers → Kafka Broker → Consumers - A Kafka Cluster = one or more brokers, coordinated by Zookeeper
📂 Part 9 — Topics, Partitions & Offsets
Topics
- Feed name to which messages are published (e.g.,
"user-events") - Producers always append to the tail (append-only log)
- Kafka prunes the "head" based on age, size, or key
- Consumers track position with an offset pointer (like an index)
Partitions
Topic: "user-events"
┌──────────────────────────────────┐
│ Partition 0: [msg0][msg1][msg2] │
│ Partition 1: [msg0][msg1][msg3] │
│ Partition 2: [msg0][msg1][msg4] │
└──────────────────────────────────┘
- Partition = ordered + immutable sequence of messages (continually appended)
- Number of partitions is configurable
- #partitions = max consumer parallelism (more partitions → more parallel consumers)
🛣️ Highway Analogy: Topic = highway; Partitions = lanes. More lanes → more cars (consumers) simultaneously.
Consumer Groups Example
Topic with 4 partitions:
Consumer Group A (2 consumers):
Consumer 1 → Partition 0, 1
Consumer 2 → Partition 2, 3
Consumer Group B (4 consumers):
Consumer 1 → Partition 0
Consumer 2 → Partition 1
Consumer 3 → Partition 2
Consumer 4 → Partition 3
⚠️ More consumers than partitions → some consumers are IDLE
Partition Offsets
- Offset = unique sequential ID per message per partition (Offset = Message ID)
- Each consumer group reads independently at its own pace
- Consumers can rewind (replay) or skip ahead by changing offset
🔁 Part 10 — Replicas
- Replicas = backup copies of partitions
- Purpose: prevent data loss ONLY — not for performance!
- Replicas are never read from, never written to directly
- They do NOT increase producer/consumer parallelism
LinkedIn uses
numReplicas = 2→ can tolerate 1 broker dying
🌊 Part 11 — Kafka Streams: KStream & KTable
KStream — The Event Log
- Represents an unbounded sequence of immutable events
- Each record = append-only event
- Records are never updated or deleted
- Supports windowing (e.g., last 5 minutes of data)
IoT Sensor events:
(sensor1, 30°C)
(sensor1, 32°C) ← KStream keeps ALL of these
(sensor1, 31°C)
KTable — The Current State
- Represents the latest value per key
- New records overwrite old ones for the same key
- Like a database table / materialized view
Same sensor events → KTable:
sensor1 → 31°C ← only the LATEST value
KStream vs KTable
| Property | KStream | KTable |
|---|---|---|
| Interpretation | INSERT / append | UPSERT / overwrite |
| Use case | Transaction history | Current balance |
| Analogy | Bank statement (every tx) | Bank balance (current state) |
Converting Between Them
KStream → KTable : group + aggregate the stream
KTable → KStream : emit changes as a changelog stream
⚙️ Part 12 — Kafka Streams Operations
Stateless Transformations (No memory of past data)
| Operation | Swift Equivalent | Description |
|---|---|---|
filter | .filter {} | Keep records that match a condition |
map | .map {} | Transform each record into a different one |
mapValues | .map { $0.value } | Transform only the value, keep key |
flatMap | .flatMap {} | Transform each record into 0 or more records |
flatMapValues | - | Transform values into 0 or more values |
Stateful Transformations (Remember past data → result is KTable)
| Operation | Description |
|---|---|
countByKey | Count occurrences per key → ever-updating KTable |
reduceByKey | Combine values with a reducer function → ever-updating KTable |
These become a KTable because the output changes over time as more data arrives.
Aggregates on KStream (Must Group First!)
Input:
(dev1, 30), (dev2, 25), (dev1, 32), (dev2, 28)
After groupByKey:
dev1 → [30, 32]
dev2 → [25, 28]
Now compute: sum, avg, max, min, count
🔧 Part 13 — Process Topology
- Processor Topology = graph of stream processors (nodes) connected by streams (edges)
- Stream = unbounded, continuously updating dataset (ordered, immutable key-value pairs)
- Source Processor = reads from Kafka topic → forwards to downstream
- Sink Processor = sends records from upstream → back to a Kafka topic
Kafka Topic → [Source Processor] → [Transform Processor] → [Sink Processor] → Kafka Topic
🛡️ Part 14 — Kafka Streams Fault Tolerance
- If a Stream Processor crashes (e.g., D crashes):
- Its partition (P4) is reassigned to another processor (C now handles P3 + P4)
- The crashed processor's changelog (state changes) is stored as a separate Kafka topic
- When a new processor takes over, it replays the changelog to restore state
- Result: no data loss, state is fully recoverable!
⚔️ Part 15 — Kafka Streams vs. Spark Streaming
| # | Spark Streaming | Kafka Streams |
|---|---|---|
| Processing unit | Micro-batches | Per record (true real-time) |
| Cluster needed? | Requires separate cluster | ❌ No cluster — it's a library |
| Scaling | Needs reconfiguration | Add Java processes |
| Semantics | At-least-once | Exactly-once |
| Best for | groupBy, ML, window functions | Row parsing, data cleansing, record-at-a-time |
| Architecture | Standalone framework | Microservice-friendly (just a library) |
| Storage | RDD (distributed, cached) | Topic (buffer memory) |
| Language | Java, Scala, R, Python | Java primary |
🌍 Part 16 — Real-World Use Cases & Examples
Top 5 Kafka Use Cases
| # | Use Case | Pipeline |
|---|---|---|
| 1 | Log Analysis | Services → Kafka → Elastic → Kibana |
| 2 | Recommendations | User clicks → Kafka → Flink → ML Models |
| 3 | System Monitoring | Metrics → Kafka → Flink → Alerts |
| 4 | Change Data Capture | Source DB → tx log → Kafka → replica DBs |
| 5 | System Migration | Old services → Kafka → New services |
Uber Kafka Architecture
Producers: Rider App, Driver App, API Services, Dispatch (GPS), Map Services
↓
KAFKA (Real-time)
↓ ↓ ↓ ↓
Surge ELK Storm Samza → Alerting, Analytics, Debugging
↓
HADOOP (Batch)
↓
VERTICA → Analytics Reporting
IoT EHR (Health) Pipeline
IoT Device (Producer)
↓
Kafka Topic: patient-vitals
↓
KStream (event processing)
↓──────────────────────────────┐
↓ ↓
Alert Stream (filter HR > 100) KTable (latest patient state)
↓ ↓
Alert System Dashboard (Consumer)
IoT Device State Example
// Convert incoming events (KStream) to current device state (KTable)
KTable<String, DeviceState> deviceTable =
stream
.groupByKey()
.reduce((oldValue, newValue) -> newValue); // keep only latest
// Input KStream:
// (dev1, 30°C, ON, 90%)
// (dev2, 28°C, ON, 80%)
// (dev1, 32°C, ON, 88%) ← overwrites
// (dev2, 29°C, OFF, 78%) ← overwrites
// Resulting KTable:
// dev1 → (32°C, ON, 88%)
// dev2 → (29°C, OFF, 78%)🪵 Part 17 — Logstash (ELK Stack)
Logstash is part of the ELK stack (Elasticsearch + Logstash + Kibana), acquired by Elastic.
Pipeline
Beats → Kafka → Logstash → Elasticsearch → Kibana
Three Stages
| Stage | Description | Examples |
|---|---|---|
| Inputs | Ingest data from any source | Beats, S3, DB, Kafka |
| Filters | Parse & transform on the fly | GROK, GeoIP, Fingerprint, Date |
| Outputs | Route data to any destination | Elasticsearch, Syslog, Slack |
What Logstash Filters Can Do
- Derive structure from unstructured logs (GROK)
- Decode geo coordinates from IP addresses
- Anonymize PII data / exclude sensitive fields
- Works independently of data source, format, or schema