Kappa-архитектура была предложена Джеем Крепсом (Jay Kreps, создатель Kafka) как упрощение Lambda-архитектуры. Основная идея: использовать один слой обработки (stream processing) вместо двух (batch + stream).
Принцип
Все данные -- это поток событий. Если нужно переобработать исторические данные, достаточно перечитать лог событий с нужной точки.
Шаг 1: Все данные записываются в иммутабельный лог
<?php
declare(strict_types=1);
/**
* Immutable event log: append-only storage
*/
final class EventLog
{
public function __construct(
private readonly KafkaProducer $producer,
) {}
/**
* Every state change is an event in the log.
* Events are never modified or deleted.
*/
public function append(string $streamId, array $event): void
{
$this->producer->send($streamId, [
'event_id' => uuid_create(),
'stream_id' => $streamId,
'type' => $event['type'],
'data' => $event['data'],
'metadata' => [
'occurred_at' => (new \DateTimeImmutable())->format('c'),
'version' => $event['version'] ?? 1,
'source' => $event['source'] ?? 'unknown',
],
]);
}
}
// Usage: all changes are events
$log = new EventLog($producer);
$log->append('user:123', [
'type' => 'user.registered',
'data' => ['email' => '[email protected]', 'name' => 'John'],
'version' => 1,
]);
$log->append('user:123', [
'type' => 'user.email_changed',
'data' => ['old_email' => '[email protected]', 'new_email' => '[email protected]'],
'version' => 2,
]);
package kappa
import (
"context"
"time"
"github.com/google/uuid"
)
// KafkaProducer is an interface for sending messages to Kafka.
type KafkaProducer interface {
Send(ctx context.Context, key string, payload map[string]any) error
}
// EventLog provides append-only immutable event storage.
type EventLog struct {
producer KafkaProducer
}
func NewEventLog(producer KafkaProducer) *EventLog {
return &EventLog{producer: producer}
}
// Append writes an event to the immutable log.
// Events are never modified or deleted.
func (l *EventLog) Append(ctx context.Context, streamID string, event map[string]any) error {
version := 1
if v, ok := event["version"].(int); ok {
version = v
}
source := "unknown"
if s, ok := event["source"].(string); ok {
source = s
}
return l.producer.Send(ctx, streamID, map[string]any{
"event_id": uuid.New().String(),
"stream_id": streamID,
"type": event["type"],
"data": event["data"],
"metadata": map[string]any{
"occurred_at": time.Now().Format(time.RFC3339),
"version": version,
"source": source,
},
})
}
using System.Text.Json;
using Confluent.Kafka;
public sealed record EventInput(
string Type,
IReadOnlyDictionary<string, object?> Data,
int Version = 1,
string Source = "unknown");
/// Immutable event log: append-only storage
public sealed class EventLog(IProducer<string, string> producer, string topic)
{
// Every state change is an event in the log.
// Events are never modified or deleted.
public async Task AppendAsync(string streamId, EventInput evt, CancellationToken ct = default)
{
var envelope = new
{
event_id = Guid.NewGuid().ToString(),
stream_id = streamId,
type = evt.Type,
data = evt.Data,
metadata = new
{
occurred_at = DateTimeOffset.UtcNow.ToString("o"),
version = evt.Version,
source = evt.Source,
},
};
await producer.ProduceAsync(
topic,
new Message<string, string>
{
Key = streamId,
Value = JsonSerializer.Serialize(envelope),
},
ct);
}
}
// Usage: all changes are events
var log = new EventLog(producer, "events");
await log.AppendAsync("user:123", new EventInput(
"user.registered",
new Dictionary<string, object?> { ["email"] = "[email protected]", ["name"] = "John" },
Version: 1));
await log.AppendAsync("user:123", new EventInput(
"user.email_changed",
new Dictionary<string, object?>
{
["old_email"] = "[email protected]",
["new_email"] = "[email protected]",
},
Version: 2));
import json
import uuid
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any
from confluent_kafka import Producer
@dataclass(frozen=True)
class EventInput:
type: str
data: dict[str, Any] = field(default_factory=dict)
version: int = 1
source: str = "unknown"
class EventLog:
"""Immutable event log: append-only storage."""
def __init__(self, producer: Producer, topic: str) -> None:
self._producer = producer
self._topic = topic
def append(self, stream_id: str, event: EventInput) -> None:
"""Every state change is an event in the log.
Events are never modified or deleted.
"""
envelope = {
"event_id": str(uuid.uuid4()),
"stream_id": stream_id,
"type": event.type,
"data": event.data,
"metadata": {
"occurred_at": datetime.now(timezone.utc).isoformat(),
"version": event.version,
"source": event.source,
},
}
self._producer.produce(
self._topic,
key=stream_id,
value=json.dumps(envelope),
)
# Usage: all changes are events
log = EventLog(producer, "events")
log.append("user:123", EventInput(
type="user.registered",
data={"email": "[email protected]", "name": "John"},
version=1,
))
log.append("user:123", EventInput(
type="user.email_changed",
data={"old_email": "[email protected]", "new_email": "[email protected]"},
version=2,
))
### Шаг 2: Stream processor создаёт производные представления
<?php
declare(strict_types=1);
/**
* Stream processor builds derived views from the event log
*/
final class UserViewProcessor
{
public function __construct(
private readonly \PDO $db,
private readonly RedisClient $redis,
) {}
/**
* Process each event and update materialized views
*/
public function process(array $event): void
{
match ($event['type']) {
'user.registered' => $this->handleRegistered($event),
'user.email_changed' => $this->handleEmailChanged($event),
'user.deactivated' => $this->handleDeactivated($event),
default => null,
};
}
private function handleRegistered(array $event): void
{
$data = $event['data'];
// Materialized view in PostgreSQL
$stmt = $this->db->prepare(<<<SQL
INSERT INTO users_view (id, email, name, status, created_at, updated_at)
VALUES (:id, :email, :name, 'active', :created_at, :created_at)
ON CONFLICT (id) DO NOTHING
SQL);
$stmt->execute([
'id' => $event['stream_id'],
'email' => $data['email'],
'name' => $data['name'],
'created_at' => $event['metadata']['occurred_at'],
]);
// Cache for fast reads
$this->redis->hMSet("user:{$event['stream_id']}", [
'email' => $data['email'],
'name' => $data['name'],
'status' => 'active',
]);
}
private function handleEmailChanged(array $event): void
{
$data = $event['data'];
$stmt = $this->db->prepare(<<<SQL
UPDATE users_view
SET email = :email, updated_at = :updated_at
WHERE id = :id
SQL);
$stmt->execute([
'id' => $event['stream_id'],
'email' => $data['new_email'],
'updated_at' => $event['metadata']['occurred_at'],
]);
$this->redis->hSet("user:{$event['stream_id']}", 'email', $data['new_email']);
}
private function handleDeactivated(array $event): void
{
$stmt = $this->db->prepare(<<<SQL
UPDATE users_view
SET status = 'inactive', updated_at = :updated_at
WHERE id = :id
SQL);
$stmt->execute([
'id' => $event['stream_id'],
'updated_at' => $event['metadata']['occurred_at'],
]);
$this->redis->hSet("user:{$event['stream_id']}", 'status', 'inactive');
}
}
using System.Text.Json;
using Npgsql;
using StackExchange.Redis;
/// Stream processor builds derived views from the event log
public sealed class UserViewProcessor(NpgsqlDataSource db, IDatabase redis)
{
// Process each event and update materialized views
public Task ProcessAsync(JsonElement evt, CancellationToken ct = default)
=> evt.GetProperty("type").GetString() switch
{
"user.registered" => HandleRegisteredAsync(evt, ct),
"user.email_changed" => HandleEmailChangedAsync(evt, ct),
"user.deactivated" => HandleDeactivatedAsync(evt, ct),
_ => Task.CompletedTask,
};
private async Task HandleRegisteredAsync(JsonElement evt, CancellationToken ct)
{
var data = evt.GetProperty("data");
var streamId = evt.GetProperty("stream_id").GetString()!;
var occurredAt = evt.GetProperty("metadata").GetProperty("occurred_at").GetDateTimeOffset();
var email = data.GetProperty("email").GetString();
var name = data.GetProperty("name").GetString();
// Materialized view in PostgreSQL
await using var cmd = db.CreateCommand(
"""
INSERT INTO users_view (id, email, name, status, created_at, updated_at)
VALUES ($1, $2, $3, 'active', $4, $4)
ON CONFLICT (id) DO NOTHING
""");
cmd.Parameters.AddWithValue(streamId);
cmd.Parameters.AddWithValue(email!);
cmd.Parameters.AddWithValue(name!);
cmd.Parameters.AddWithValue(occurredAt);
await cmd.ExecuteNonQueryAsync(ct);
// Cache for fast reads
await redis.HashSetAsync($"user:{streamId}", new HashEntry[]
{
new("email", email),
new("name", name),
new("status", "active"),
});
}
private async Task HandleEmailChangedAsync(JsonElement evt, CancellationToken ct)
{
var streamId = evt.GetProperty("stream_id").GetString()!;
var occurredAt = evt.GetProperty("metadata").GetProperty("occurred_at").GetDateTimeOffset();
var newEmail = evt.GetProperty("data").GetProperty("new_email").GetString();
await using var cmd = db.CreateCommand(
"UPDATE users_view SET email = $1, updated_at = $2 WHERE id = $3");
cmd.Parameters.AddWithValue(newEmail!);
cmd.Parameters.AddWithValue(occurredAt);
cmd.Parameters.AddWithValue(streamId);
await cmd.ExecuteNonQueryAsync(ct);
await redis.HashSetAsync($"user:{streamId}", "email", newEmail);
}
private async Task HandleDeactivatedAsync(JsonElement evt, CancellationToken ct)
{
var streamId = evt.GetProperty("stream_id").GetString()!;
var occurredAt = evt.GetProperty("metadata").GetProperty("occurred_at").GetDateTimeOffset();
await using var cmd = db.CreateCommand(
"UPDATE users_view SET status = 'inactive', updated_at = $1 WHERE id = $2");
cmd.Parameters.AddWithValue(occurredAt);
cmd.Parameters.AddWithValue(streamId);
await cmd.ExecuteNonQueryAsync(ct);
await redis.HashSetAsync($"user:{streamId}", "status", "inactive");
}
}
from typing import Any
import asyncpg
import redis.asyncio as aioredis
class UserViewProcessor:
"""Stream processor builds derived views from the event log."""
def __init__(self, pool: asyncpg.Pool, redis: aioredis.Redis) -> None:
self._pool = pool
self._redis = redis
async def process(self, event: dict[str, Any]) -> None:
"""Process each event and update materialized views."""
handlers = {
"user.registered": self._handle_registered,
"user.email_changed": self._handle_email_changed,
"user.deactivated": self._handle_deactivated,
}
handler = handlers.get(event["type"])
if handler is not None:
await handler(event)
async def _handle_registered(self, event: dict[str, Any]) -> None:
data = event["data"]
occurred_at = event["metadata"]["occurred_at"]
# Materialized view in PostgreSQL
await self._pool.execute(
"""
INSERT INTO users_view (id, email, name, status, created_at, updated_at)
VALUES ($1, $2, $3, 'active', $4, $4)
ON CONFLICT (id) DO NOTHING
""",
event["stream_id"], data["email"], data["name"], occurred_at,
)
# Cache for fast reads
await self._redis.hset(
f"user:{event['stream_id']}",
mapping={"email": data["email"], "name": data["name"], "status": "active"},
)
async def _handle_email_changed(self, event: dict[str, Any]) -> None:
data = event["data"]
await self._pool.execute(
"UPDATE users_view SET email = $1, updated_at = $2 WHERE id = $3",
data["new_email"], event["metadata"]["occurred_at"], event["stream_id"],
)
await self._redis.hset(f"user:{event['stream_id']}", "email", data["new_email"])
async def _handle_deactivated(self, event: dict[str, Any]) -> None:
await self._pool.execute(
"UPDATE users_view SET status = 'inactive', updated_at = $1 WHERE id = $2",
event["metadata"]["occurred_at"], event["stream_id"],
)
await self._redis.hset(f"user:{event['stream_id']}", "status", "inactive")
### Шаг 3: При изменении логики -- переобработка (Replay)
<?php
declare(strict_types=1);
/**
* Replay: reprocess all events when logic changes
*/
final class StreamReplayManager
{
public function __construct(
private readonly KafkaConsumer $consumer,
private readonly \PDO $db,
) {}
/**
* Steps to deploy new processing logic:
*
* 1. Deploy new version of processor (v2)
* 2. v2 reads from the beginning of the log
* 3. v2 writes to new output tables/views
* 4. When v2 catches up with real-time, switch reads to v2
* 5. Decommission v1
*/
public function replay(
string $topic,
string $newConsumerGroup,
callable $newProcessor,
): void {
// Create new consumer group, starting from the beginning
$replayConsumer = new KafkaConsumer(
brokers: 'kafka1:9092',
groupId: $newConsumerGroup, // e.g., 'user-view-v2'
topics: [$topic],
);
// Track replay progress
$processedCount = 0;
$startTime = microtime(true);
$replayConsumer->run(function (array $event) use ($newProcessor, &$processedCount, $startTime): void {
$newProcessor($event);
$processedCount++;
if ($processedCount % 10000 === 0) {
$elapsed = microtime(true) - $startTime;
$rate = $processedCount / $elapsed;
error_log("Replay progress: {$processedCount} events, {$rate:.0f} events/sec");
}
});
}
}
package kappa
import (
"context"
"fmt"
"log"
"time"
)
// StreamReplayManager replays all events when processing logic changes.
type StreamReplayManager struct{}
// Replay creates a new consumer group from the beginning and reprocesses all events.
//
// Steps:
// 1. Deploy new processor (v2)
// 2. v2 reads from the beginning of the log
// 3. v2 writes to new output tables/views
// 4. When v2 catches up, switch reads to v2
// 5. Decommission v1
func (m *StreamReplayManager) Replay(
ctx context.Context,
brokers, topic, newConsumerGroup string,
processor func(ctx context.Context, event map[string]any) error,
) error {
consumer, err := NewConsumer(brokers, newConsumerGroup, []string{topic})
if err != nil {
return fmt.Errorf("create replay consumer: %w", err)
}
defer consumer.Close()
processedCount := 0
startTime := time.Now()
return consumer.Run(ctx, func(ctx context.Context, payload map[string]any, key string) error {
if err := processor(ctx, payload); err != nil {
return err
}
processedCount++
if processedCount%10000 == 0 {
elapsed := time.Since(startTime).Seconds()
rate := float64(processedCount) / elapsed
log.Printf("Replay progress: %d events, %.0f events/sec", processedCount, rate)
}
return nil
})
}
using System.Diagnostics;
using System.Text.Json;
using Confluent.Kafka;
using Microsoft.Extensions.Logging;
/// Replay: reprocess all events when logic changes
public sealed class StreamReplayManager(ILogger<StreamReplayManager> logger)
{
/// Steps to deploy new processing logic:
///
/// 1. Deploy new version of processor (v2)
/// 2. v2 reads from the beginning of the log
/// 3. v2 writes to new output tables/views
/// 4. When v2 catches up with real-time, switch reads to v2
/// 5. Decommission v1
public async Task ReplayAsync(
string brokers,
string topic,
string newConsumerGroup,
Func<JsonElement, CancellationToken, Task> newProcessor,
CancellationToken ct = default)
{
// New consumer group + AutoOffsetReset.Earliest means reading from the beginning
var config = new ConsumerConfig
{
BootstrapServers = brokers,
GroupId = newConsumerGroup, // e.g., "user-view-v2"
AutoOffsetReset = AutoOffsetReset.Earliest,
EnableAutoCommit = false,
};
using var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe(topic);
// Track replay progress
var processedCount = 0;
var stopwatch = Stopwatch.StartNew();
while (!ct.IsCancellationRequested)
{
var result = consumer.Consume(ct);
var evt = JsonDocument.Parse(result.Message.Value).RootElement;
await newProcessor(evt, ct);
consumer.Commit(result);
processedCount++;
if (processedCount % 10_000 == 0)
{
var rate = processedCount / stopwatch.Elapsed.TotalSeconds;
logger.LogInformation(
"Replay progress: {Count} events, {Rate:F0} events/sec", processedCount, rate);
}
}
}
}
import json
import logging
import time
from typing import Any, Awaitable, Callable
from confluent_kafka import Consumer
logger = logging.getLogger(__name__)
Processor = Callable[[dict[str, Any]], Awaitable[None]]
class StreamReplayManager:
"""Replay: reprocess all events when logic changes."""
async def replay(
self,
brokers: str,
topic: str,
new_consumer_group: str,
new_processor: Processor,
) -> None:
"""Steps to deploy new processing logic:
1. Deploy new version of processor (v2)
2. v2 reads from the beginning of the log
3. v2 writes to new output tables/views
4. When v2 catches up with real-time, switch reads to v2
5. Decommission v1
"""
# New consumer group + earliest offset means reading from the beginning
consumer = Consumer({
"bootstrap.servers": brokers,
"group.id": new_consumer_group, # e.g., "user-view-v2"
"auto.offset.reset": "earliest",
"enable.auto.commit": False,
})
consumer.subscribe([topic])
# Track replay progress
processed_count = 0
start_time = time.monotonic()
try:
while True:
message = consumer.poll(1.0)
if message is None:
continue
if message.error() is not None:
raise RuntimeError(f"replay consume failed: {message.error()}")
await new_processor(json.loads(message.value()))
consumer.commit(message)
processed_count += 1
if processed_count % 10_000 == 0:
elapsed = time.monotonic() - start_time
rate = processed_count / elapsed
logger.info(
"Replay progress: %d events, %.0f events/sec", processed_count, rate
)
finally:
consumer.close()
## Kappa vs Lambda: подробное сравнение
Аспект
Lambda
Kappa
Количество путей обработки
2 (batch + speed)
1 (stream only)
Кодовые базы
2 разных
1 единая
Переобработка данных
Batch layer пересчитывает
Replay event log
Сложность эксплуатации
Высокая
Средняя
Требования к хранению
Меньше (агрегаты)
Больше (весь лог)
Точность
Batch гарантирует
Зависит от семантики
Время переобработки
Часы (batch job)
Зависит от объёма лога
Технологический стек
Hadoop + Kafka/Flink
Kafka (или Flink)
Когда выбрать Lambda
Batch и stream логика принципиально различается
Нужны сложные ML-модели, обучаемые на всех данных
Объём исторических данных слишком велик для replay
Команда уже имеет инфраструктуру для batch (Hadoop, Spark)
Когда выбрать Kappa
Одна логика обработки для всех данных
Event-driven микросервисы
Данные естественно представлены как поток событий
Нужна простота эксплуатации
Допустимо хранить полный лог событий
Stream Processing Engine
Kappa-архитектура требует мощного stream processing engine. Основные варианты:
Инструмент
Описание
PHP-совместимость
Kafka Streams
Библиотека для JVM
Нет (Java only)
Apache Flink
Распределённый движок
Нет (Java/Python)
Apache Spark Streaming
Микро-батчи
Нет (Scala/Python)
PHP + rdkafka
Consumer loop в PHP
Да
PHP как Stream Processor
PHP вполне подходит для stream processing на уровне consumer applications:
<?php
declare(strict_types=1);
/**
* PHP stream processor with windowed aggregation
*/
final class WindowedAggregator
{
/** @var array<string, array{count: int, sum: float, window_start: int}> */
private array $windows = [];
public function __construct(
private readonly RedisClient $redis,
private readonly int $windowSizeSeconds = 60,
) {}
/**
* Tumbling window aggregation
*/
public function aggregate(array $event): void
{
$eventTime = strtotime($event['metadata']['occurred_at']);
$windowStart = $eventTime - ($eventTime % $this->windowSizeSeconds);
$windowKey = "window:{$windowStart}:{$event['data']['category']}";
if (!isset($this->windows[$windowKey])) {
$this->windows[$windowKey] = [
'count' => 0,
'sum' => 0.0,
'window_start' => $windowStart,
];
}
$this->windows[$windowKey]['count']++;
$this->windows[$windowKey]['sum'] += (float) $event['data']['amount'];
// Flush completed windows
$this->flushCompletedWindows(time());
}
private function flushCompletedWindows(int $currentTime): void
{
foreach ($this->windows as $key => $window) {
$windowEnd = $window['window_start'] + $this->windowSizeSeconds;
if ($currentTime > $windowEnd + 10) { // 10 sec grace period for late events
// Emit window result
$this->redis->hMSet("result:{$key}", [
'count' => $window['count'],
'sum' => $window['sum'],
'avg' => $window['sum'] / $window['count'],
]);
unset($this->windows[$key]);
}
}
}
}
package kappa
import (
"context"
"fmt"
"sync"
"time"
"github.com/redis/go-redis/v9"
)
// WindowData holds aggregation state for a single window.
type WindowData struct {
Count int
Sum float64
WindowStart int64
}
// WindowedAggregator performs tumbling window aggregation over a stream.
type WindowedAggregator struct {
redis *redis.Client
windowSizeSeconds int64
mu sync.Mutex
windows map[string]*WindowData
}
func NewWindowedAggregator(rdb *redis.Client, windowSize int64) *WindowedAggregator {
return &WindowedAggregator{
redis: rdb,
windowSizeSeconds: windowSize,
windows: make(map[string]*WindowData),
}
}
// Aggregate adds an event to the correct tumbling window.
func (a *WindowedAggregator) Aggregate(ctx context.Context, event map[string]any) error {
meta, _ := event["metadata"].(map[string]any)
data, _ := event["data"].(map[string]any)
occurredAt, _ := meta["occurred_at"].(string)
eventTime, _ := time.Parse(time.RFC3339, occurredAt)
unixTime := eventTime.Unix()
windowStart := unixTime - (unixTime % a.windowSizeSeconds)
category, _ := data["category"].(string)
windowKey := fmt.Sprintf("window:%d:%s", windowStart, category)
a.mu.Lock()
if _, ok := a.windows[windowKey]; !ok {
a.windows[windowKey] = &WindowData{WindowStart: windowStart}
}
a.windows[windowKey].Count++
amount, _ := data["amount"].(float64)
a.windows[windowKey].Sum += amount
a.mu.Unlock()
return a.flushCompletedWindows(ctx, time.Now().Unix())
}
func (a *WindowedAggregator) flushCompletedWindows(ctx context.Context, currentTime int64) error {
a.mu.Lock()
defer a.mu.Unlock()
for key, w := range a.windows {
windowEnd := w.WindowStart + a.windowSizeSeconds
if currentTime > windowEnd+10 { // 10 sec grace period
avg := w.Sum / float64(w.Count)
a.redis.HSet(ctx, "result:"+key, map[string]any{
"count": w.Count, "sum": w.Sum, "avg": avg,
})
delete(a.windows, key)
}
}
return nil
}
## Управление состоянием (State Management)
В Kappa-архитектуре stream processor часто хранит состояние:
<?php
declare(strict_types=1);
/**
* Stateful stream processor with Redis-backed state
*/
final class StatefulProcessor
{
public function __construct(
private readonly RedisClient $redis,
) {}
/**
* Count unique users per day using HyperLogLog
*/
public function trackUniqueUser(array $event): void
{
$date = date('Y-m-d', strtotime($event['metadata']['occurred_at']));
$userId = $event['data']['user_id'];
// HyperLogLog: memory-efficient unique counting
$this->redis->pfAdd("unique_users:{$date}", [$userId]);
}
/**
* Detect duplicate events using idempotency key
*/
public function isDuplicate(array $event): bool
{
$eventId = $event['event_id'];
// SET NX with TTL for deduplication window
$isNew = $this->redis->set(
"processed:{$eventId}",
'1',
['NX', 'EX' => 86400], // 24 hour dedup window
);
return !$isNew;
}
/**
* Running average with exponential decay
*/
public function updateRunningAverage(string $metric, float $value, float $alpha = 0.1): float
{
$currentAvg = (float) $this->redis->get("avg:{$metric}");
// Exponential moving average
$newAvg = $currentAvg === 0.0
? $value
: $alpha * $value + (1 - $alpha) * $currentAvg;
$this->redis->set("avg:{$metric}", (string) $newAvg);
return $newAvg;
}
}