Chapter 11 - Akka

Updated 4 Oct 2026

🧠 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:

ModelUnitResponseBlocking?Sync?Addresses Recipient?
DatabasesDataNoneNon-blockingAsync❌ No
Services (RPC)Function callYesBlockingSync✅ Yes
Message-PassingMessagesMaybeNon-blockingAsync✅ Yes

🍵 Analogy

ModelReal Life
DatabaseWhiteboard in the office — anyone reads/writes, no direct recipient
Service (RPC)Phone call — you wait until the other person picks up and responds
Message-PassingLINE 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.

Actor=State+Behavior+Mailbox\text{Actor} = \text{State} + \text{Behavior} + \text{Mailbox}

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?

  1. Send messages to other actors (non-blocking)
  2. Create new child actors (spawn!)
  3. 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-OrientedActorTask-Oriented
Core UnitObjectActorTask
CommunicationMethod callsAsync messagesScheduled tasks
Key BenefitSeparation of concernsOO + asyncAsync/parallel decoupling

Advantages vs. Drawbacks

✅ Advantages❌ Drawbacks
No explicit locks → no deadlocksCreating many actors → performance/latency cost
Fault-tolerant "let it crash" philosophyNeed to track where actors live (location overhead)
Dynamic parallelism via spawningActive 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

ModulePurpose
Akka ActorsCore: concurrency + distribution
Akka ClusterResilient distribution across multiple nodes
Akka StreamsAsync, non-blocking, backpressured reactive streams
Akka HttpStreaming-first HTTP server/client
Cluster ShardingReference actors by identity, not location
Akka PersistencePersist actor state for fault tolerance (≠ ephemeral!)
Distributed DataEventually consistent key-value store
AlpakkaConnectors to other tech stacks

Akka Actor Internals

Actor=State+Behavior+Mailbox\text{Actor} = \text{State} + \text{Behavior} + \text{Mailbox}

  • 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 PartMeaning
AbstractActorInherit actor behavior, state, mailbox
ReceivePattern 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)

ActionWhat Happens
restartKill and restart the child actor (fresh state)
resumeIgnore the failure, continue processing
stopPermanently stop the child actor
escalatePass 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

GuaranteeDeliveryDuplicates?PerformanceHow
at-most-once0 or 1 timesNever🟢 FastestFire-and-forget
at-least-once1+ timesPossible🟡 OKSend + acknowledge
exactly-onceExactly 1 timeNever🔴 SlowestSend + 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

ModelHowRisk
PushProducer immediately sends work to consumers (queued in consumer's inbox)⚠️ Congestion / message drops if consumer is slow
PullConsumer 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?

CompanyUsage
LinkedIn400 nodes, 18k topics, 220B msg/day (peak 3.2M msg/s)
NetflixReal-time monitoring and event processing
TwitterPart of Storm real-time pipelines
SpotifyLog delivery: reduced from 4h → 10 seconds
UberGPS logs, surge pricing, driver/rider events

🧱 Part 8 — Kafka Core Concepts

ConceptDescription
BrokerA single Kafka node (server) in the cluster
TopicA named stream/category of records (like a channel)
ProducerPushes messages into a topic
ConsumerPulls messages from a topic
Data RetentionMessages kept based on time or size limit
ZookeeperStores 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

Consumer Position=(offset, partition, topic)\text{Consumer Position} = (\text{offset},\ \text{partition},\ \text{topic})

  • 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

Brokers tolerated=numReplicas−1\text{Brokers tolerated} = \text{numReplicas} - 1

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

PropertyKStreamKTable
InterpretationINSERT / appendUPSERT / overwrite
Use caseTransaction historyCurrent balance
AnalogyBank 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)

OperationSwift EquivalentDescription
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)

OperationDescription
countByKeyCount occurrences per key → ever-updating KTable
reduceByKeyCombine 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

  1. If a Stream Processor crashes (e.g., D crashes):
    • Its partition (P4) is reassigned to another processor (C now handles P3 + P4)
  2. The crashed processor's changelog (state changes) is stored as a separate Kafka topic
  3. When a new processor takes over, it replays the changelog to restore state
  4. Result: no data loss, state is fully recoverable!

⚔️ Part 15 — Kafka Streams vs. Spark Streaming

#Spark StreamingKafka Streams
Processing unitMicro-batchesPer record (true real-time)
Cluster needed?Requires separate cluster❌ No cluster — it's a library
ScalingNeeds reconfigurationAdd Java processes
SemanticsAt-least-onceExactly-once
Best forgroupBy, ML, window functionsRow parsing, data cleansing, record-at-a-time
ArchitectureStandalone frameworkMicroservice-friendly (just a library)
StorageRDD (distributed, cached)Topic (buffer memory)
LanguageJava, Scala, R, PythonJava primary

🌍 Part 16 — Real-World Use Cases & Examples

Top 5 Kafka Use Cases

#Use CasePipeline
1Log AnalysisServices → Kafka → Elastic → Kibana
2RecommendationsUser clicks → Kafka → Flink → ML Models
3System MonitoringMetrics → Kafka → Flink → Alerts
4Change Data CaptureSource DB → tx log → Kafka → replica DBs
5System MigrationOld 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

StageDescriptionExamples
InputsIngest data from any sourceBeats, S3, DB, Kafka
FiltersParse & transform on the flyGROK, GeoIP, Fingerprint, Date
OutputsRoute data to any destinationElasticsearch, 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