MidТеория25 min

Kappa архитектура

Отличия от Lambda, когда использовать, stream processing как единый слой обработки данных

Идея Kappa

Kappa-архитектура была предложена Джеем Крепсом (Jay Kreps, создатель Kafka) как упрощение Lambda-архитектуры. Основная идея: использовать один слой обработки (stream processing) вместо двух (batch + stream).

Принцип

Все данные -- это поток событий. Если нужно переобработать исторические данные, достаточно перечитать лог событий с нужной точки.

Lambda:                          Kappa:

┌───────────┐                    ┌───────────┐
│   Batch   │                    │           │
│  Layer    │──┐                 │  Stream   │
└───────────┘  │  ┌─────────┐   │Processing │──── Serving
               ├──│ Serving │   │  (единый) │     Layer
┌───────────┐  │  └─────────┘   │           │
│  Speed    │──┘                 └─────┬─────┘
│  Layer    │                          │
└───────────┘                    ┌─────┴─────┐
                                 │Immutable  │
                                 │   Log     │
                                 └───────────┘

Как работает Kappa

Шаг 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,
]);
### Шаг 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');
    }
}
### Шаг 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");
            }
        });
    }
}
## 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]);
            }
        }
    }
}
## Управление состоянием (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;
    }
}
## Практический пример: аналитика заказов
<?php

declare(strict_types=1);

/**
 * Complete Kappa pipeline for order analytics
 */
final class OrderAnalyticsPipeline
{
    public function __construct(
        private readonly \PDO $db,
        private readonly RedisClient $redis,
        private readonly KafkaProducer $outputProducer,
    ) {}

    /**
     * Single stream processor handles all analytics
     */
    public function process(array $event): void
    {
        if ($this->isDuplicate($event['event_id'])) {
            return;
        }

        match ($event['type']) {
            'order.created' => $this->onOrderCreated($event['data']),
            'order.completed' => $this->onOrderCompleted($event['data']),
            'order.refunded' => $this->onOrderRefunded($event['data']),
            default => null,
        };
    }

    private function onOrderCreated(array $data): void
    {
        // Real-time counters
        $this->redis->incr('stats:orders:total');
        $this->redis->incrByFloat('stats:revenue:pending', $data['total']);

        // Per-category stats
        foreach ($data['items'] as $item) {
            $this->redis->hIncrBy('stats:category:count', $item['category'], 1);
            $this->redis->hIncrByFloat('stats:category:revenue', $item['category'], $item['price']);
        }

        // Materialized view for queries
        $stmt = $this->db->prepare(<<<SQL
            INSERT INTO order_analytics (order_id, user_id, total, category, status, created_at)
            VALUES (:order_id, :user_id, :total, :category, 'created', :created_at)
        SQL);

        $stmt->execute([
            'order_id' => $data['order_id'],
            'user_id' => $data['user_id'],
            'total' => $data['total'],
            'category' => $data['items'][0]['category'] ?? 'unknown',
            'created_at' => $data['created_at'],
        ]);
    }

    private function onOrderCompleted(array $data): void
    {
        $this->redis->incrByFloat('stats:revenue:pending', -$data['total']);
        $this->redis->incrByFloat('stats:revenue:completed', $data['total']);

        $stmt = $this->db->prepare(<<<SQL
            UPDATE order_analytics SET status = 'completed' WHERE order_id = :order_id
        SQL);
        $stmt->execute(['order_id' => $data['order_id']]);
    }

    private function onOrderRefunded(array $data): void
    {
        $this->redis->incrByFloat('stats:revenue:completed', -$data['refund_amount']);
        $this->redis->incrByFloat('stats:revenue:refunded', $data['refund_amount']);

        $stmt = $this->db->prepare(<<<SQL
            UPDATE order_analytics SET status = 'refunded' WHERE order_id = :order_id
        SQL);
        $stmt->execute(['order_id' => $data['order_id']]);
    }

    private function isDuplicate(string $eventId): bool
    {
        return !$this->redis->set("dedup:{$eventId}", '1', ['NX', 'EX' => 86400]);
    }
}
## Итоги
  • Kappa упрощает Lambda, убирая batch layer в пользу единого stream processing
  • Переобработка данных выполняется через replay лога событий
  • PHP + rdkafka хорошо подходит для consumer-уровня stream processing
  • Управление состоянием через Redis или другое быстрое хранилище
  • Kappa идеальна для event-driven микросервисов с единой логикой обработки