DE Wikiarchitectures / kappa

Kappa Architecture

Kappa architecture simplifies Lambda architecture by using a single streaming pipeline for both real-time and batch processing. Instead of maintaining separate batch and speed layers, Kappa treats all data as a stream and reprocesses historical data by replaying the stream from the beginning.

Architecture Overview

Kappa architecture has two main components:

  1. Immutable Log - A durable, append-only log (Apache Kafka) that stores the complete event history. Data is never deleted or mutated.
  2. Stream Processor - A single stream processing engine (Apache Flink, Kafka Streams, Warp) that computes views from the log. The same code handles both real-time and reprocessing.

How It Works

In Kappa architecture, there is no distinction between batch and streaming code. The stream processor always reads from the log. For real-time processing, it reads from the tail of the log (latest offsets). For batch/reprocessing, it reads from the head of the log (offset 0 or a specific checkpoint).

// Kafka Streams: Kappa architecture pipeline

import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;

StreamsBuilder builder = new StreamsBuilder();

KStream<String, String> events = builder.stream("raw-events");

// Single pipeline for both real-time and reprocessing
KTable<Windowed<String>, Long> hourlyCounts = events
    .groupBy((key, value) -> extractEventType(value))
    .windowedBy(TimeWindows.of(Duration.ofHours(1)))
    .count();

// State store allows reprocessing from any point
hourlyCounts.toStream().to("hourly-event-counts");

KafkaStreams streams = new KafkaStreams(builder.build(), config);

// For real-time: start from latest
config.put(StreamsConfig.AUTO_OFFSET_RESET_CONFIG, "latest");

// For reprocessing: start from beginning
config.put(StreamsConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

// Or reprocess a specific time range using Kafka's offset management
streams.cleanUp(); // clear local state
streams.start();   // re-reads from beginning

Flink Kappa Example

-- Flink SQL: Kappa architecture - single pipeline CREATE TABLE events ( event_id   STRING, user_id    INT, event_type STRING, value      DOUBLE, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'events', 'properties.bootstrap.servers' = 'kafka:9092', 'properties.group.id' = 'kappa-pipeline', 'format' = 'json', -- Offsets managed by Flink's checkpointing 'scan.startup.mode' = 'group-offsets' ); -- Same query handles both streaming and batch INSERT INTO metrics SELECT event_type, TUMBLE_START(event_time, INTERVAL '1' HOUR) AS window_start, COUNT(*) AS event_count, SUM(value) AS total_value FROM events GROUP BY event_type, TUMBLE(event_time, INTERVAL '1' HOUR); -- To reprocess: reset offsets to earliest in Kafka -- ALTER TABLE events RESET ('scan.startup.mode' = 'earliest-offset'); -- The same pipeline re-executes from the beginning

Advantages over Lambda

  • Single codebase - No synchronization between batch and streaming logic
  • Simpler operations - One pipeline to deploy, monitor, and debug
  • Consistent results - No approximation vs accuracy gap; reprocessing produces identical output
  • Natural reprocessing - Replay the log from any point to regenerate views

Challenges

  • Stream retention - Kafka must retain the full event history, which can be expensive at high volumes
  • State management - Long-running stateful operations need efficient checkpointing and recovery
  • Windowing - Reprocessing large windows requires replaying the full event stream

When to Use Kappa

  • Your stream processing engine can handle both real-time and historical data
  • You want a single codebase for all data processing
  • Kafka retention costs are acceptable for your data volume
  • You need exactly-once reprocessing guarantees

Resources

Kappa Architecture

Kappa architecture simplifies Lambda architecture by using a single streaming pipeline for both real-time and batch processing. Instead of maintaining separate batch and speed layers, Kappa treats all data as a stream and reprocesses historical data by replaying the stream from the beginning.

Architecture Overview

Kappa architecture has two main components:

  1. Immutable Log - A durable, append-only log (Apache Kafka) that stores the complete event history. Data is never deleted or mutated.
  2. Stream Processor - A single stream processing engine (Apache Flink, Kafka Streams, Warp) that computes views from the log. The same code handles both real-time and reprocessing.

How It Works

In Kappa architecture, there is no distinction between batch and streaming code. The stream processor always reads from the log. For real-time processing, it reads from the tail of the log (latest offsets). For batch/reprocessing, it reads from the head of the log (offset 0 or a specific checkpoint).

// Kafka Streams: Kappa architecture pipeline

import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;

StreamsBuilder builder = new StreamsBuilder();

KStream<String, String> events = builder.stream("raw-events");

// Single pipeline for both real-time and reprocessing
KTable<Windowed<String>, Long> hourlyCounts = events
    .groupBy((key, value) -> extractEventType(value))
    .windowedBy(TimeWindows.of(Duration.ofHours(1)))
    .count();

// State store allows reprocessing from any point
hourlyCounts.toStream().to("hourly-event-counts");

KafkaStreams streams = new KafkaStreams(builder.build(), config);

// For real-time: start from latest
config.put(StreamsConfig.AUTO_OFFSET_RESET_CONFIG, "latest");

// For reprocessing: start from beginning
config.put(StreamsConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

// Or reprocess a specific time range using Kafka's offset management
streams.cleanUp(); // clear local state
streams.start();   // re-reads from beginning

Flink Kappa Example

-- Flink SQL: Kappa architecture - single pipeline CREATE TABLE events ( event_id   STRING, user_id    INT, event_type STRING, value      DOUBLE, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'events', 'properties.bootstrap.servers' = 'kafka:9092', 'properties.group.id' = 'kappa-pipeline', 'format' = 'json', -- Offsets managed by Flink's checkpointing 'scan.startup.mode' = 'group-offsets' ); -- Same query handles both streaming and batch INSERT INTO metrics SELECT event_type, TUMBLE_START(event_time, INTERVAL '1' HOUR) AS window_start, COUNT(*) AS event_count, SUM(value) AS total_value FROM events GROUP BY event_type, TUMBLE(event_time, INTERVAL '1' HOUR); -- To reprocess: reset offsets to earliest in Kafka -- ALTER TABLE events RESET ('scan.startup.mode' = 'earliest-offset'); -- The same pipeline re-executes from the beginning

Advantages over Lambda

  • Single codebase - No synchronization between batch and streaming logic
  • Simpler operations - One pipeline to deploy, monitor, and debug
  • Consistent results - No approximation vs accuracy gap; reprocessing produces identical output
  • Natural reprocessing - Replay the log from any point to regenerate views

Challenges

  • Stream retention - Kafka must retain the full event history, which can be expensive at high volumes
  • State management - Long-running stateful operations need efficient checkpointing and recovery
  • Windowing - Reprocessing large windows requires replaying the full event stream

When to Use Kappa

  • Your stream processing engine can handle both real-time and historical data
  • You want a single codebase for all data processing
  • Kafka retention costs are acceptable for your data volume
  • You need exactly-once reprocessing guarantees

Resources