MidТеория18 min

Apache Kafka

Topics, partitions, consumer groups, offsets, retention -- архитектура Kafka и работа через php-rdkafka

Что такое Kafka

Apache Kafka -- распределённая платформа потоковой передачи данных. Kafka работает как распределённый коммит-лог: записывает события в упорядоченный, неизменяемый журнал.

Основные свойства

  • Высокая пропускная способность -- миллионы сообщений в секунду
  • Долговременное хранение -- данные хранятся на диске, не только в памяти
  • Горизонтальная масштабируемость -- добавление брокеров без простоя
  • Отказоустойчивость -- репликация между брокерами
  • Гарантия порядка -- внутри одной партиции

Архитектура Kafka

Ключевые компоненты

┌─────────────────────────────────────────────────┐
│                  Kafka Cluster                  │
│                                                 │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐      │
│  │ Broker 1 │  │ Broker 2 │  │ Broker 3 │      │
│  │          │  │          │  │          │      │
│  │ Topic A  │  │ Topic A  │  │ Topic A  │      │
│  │ Part 0   │  │ Part 1   │  │ Part 2   │      │
│  │ (leader) │  │ (leader) │  │ (leader) │      │
│  │ Part 1   │  │ Part 2   │  │ Part 0   │      │
│  │ (replica)│  │ (replica)│  │ (replica)│      │
│  └──────────┘  └──────────┘  └──────────┘      │
│                                                 │
│  ┌──────────────────────────────────────┐       │
│  │        ZooKeeper / KRaft             │       │
│  │   (metadata, leader election)        │       │
│  └──────────────────────────────────────┘       │
└─────────────────────────────────────────────────┘

Брокер (Broker)

Сервер Kafka, который хранит данные и обслуживает клиентов. Кластер состоит из нескольких брокеров.

Топик (Topic)

Логическая категория сообщений. Аналог таблицы в БД или очереди в RabbitMQ.

Партиция (Partition)

Каждый топик разделён на партиции -- упорядоченные, неизменяемые последовательности сообщений.

Topic: orders (3 partitions)

Partition 0: [msg0] [msg3] [msg6] [msg9]  ...
Partition 1: [msg1] [msg4] [msg7] [msg10] ...
Partition 2: [msg2] [msg5] [msg8] [msg11] ...

Ключевые свойства партиций:

  • Порядок гарантирован только внутри одной партиции
  • Каждая партиция имеет одного лидера (leader) и N реплик (replicas)
  • Партиция -- единица параллелизма

Распределение по партициям

<?php

declare(strict_types=1);

/**
 * Message partitioning strategies
 */
final class PartitioningExample
{
    /**
     * Messages with the same key always go to the same partition.
     * This guarantees ordering for related events.
     */
    public function determinePartition(string $key, int $partitionCount): int
    {
        // Default Kafka partitioner: murmur2 hash
        return abs(crc32($key)) % $partitionCount;
    }

    /**
     * Example: all events for user_123 go to partition 2
     * This means all events for this user are strictly ordered
     */
    public function example(): void
    {
        $userId = 'user_123';
        $partitions = 6;
        $target = $this->determinePartition($userId, $partitions);
        // user_123 -> always partition 2 (deterministic)
    }
}
> **Правило:** сообщения с одинаковым ключом попадают в одну партицию. Используйте ключ для группировки связанных событий (все заказы одного пользователя, все события одной сессии).

Consumer Groups

Consumer group -- группа потребителей, которые совместно читают топик. Каждая партиция назначается ровно одному потребителю в группе.

Topic: orders (4 partitions)

Consumer Group "order-service":
  Consumer A: reads Part 0, Part 1
  Consumer B: reads Part 2, Part 3

Consumer Group "analytics":
  Consumer C: reads Part 0, Part 1, Part 2, Part 3

Правила назначения

Потребителей Партиций Результат
2 4 По 2 партиции на потребителя
4 4 По 1 партиции на потребителя
6 4 4 активных, 2 простаивают

Важно: количество потребителей в группе не должно превышать количество партиций. Лишние потребители будут простаивать.

Offsets

Offset -- позиция сообщения в партиции. Каждая consumer group хранит свой текущий offset для каждой партиции.

Partition 0: [0] [1] [2] [3] [4] [5] [6] [7] [8]
                              ↑              ↑
                        committed         latest
                         offset           offset

Lag = latest offset - committed offset = 8 - 4 = 4 messages

Стратегии коммита offset

Стратегия Описание Риск
Auto-commit Периодический коммит (по умолчанию) Потеря/дублирование
Manual sync Коммит после обработки Блокирует поток
Manual async Асинхронный коммит Потеря коммита

Retention (Хранение)

Kafka хранит данные на диске в течение заданного периода или до достижения лимита размера.

Параметр Описание По умолчанию
retention.ms Время хранения 7 дней
retention.bytes Макс. размер партиции Без лимита
cleanup.policy delete или compact delete

Log Compaction

При cleanup.policy=compact Kafka хранит только последнее значение для каждого ключа:

Before compaction:
[key=A, v=1] [key=B, v=1] [key=A, v=2] [key=B, v=2] [key=A, v=3]

After compaction:
[key=B, v=2] [key=A, v=3]

Полезно для хранения текущего состояния (changelog topics).

PHP + Kafka: php-rdkafka

Установка

# Install librdkafka
apt-get install librdkafka-dev

# Install PHP extension
pecl install rdkafka

Producer: отправка сообщений

<?php

declare(strict_types=1);

/**
 * Kafka producer using php-rdkafka
 */
final class KafkaProducer
{
    private \RdKafka\Producer $producer;
    private \RdKafka\ProducerTopic $topic;

    public function __construct(string $brokers, string $topicName)
    {
        $config = new \RdKafka\Conf();
        $config->set('metadata.broker.list', $brokers);
        $config->set('acks', 'all'); // Wait for all replicas
        $config->set('retries', '3');
        $config->set('retry.backoff.ms', '100');
        $config->set('enable.idempotence', 'true'); // Exactly-once producer

        // Delivery report callback
        $config->setDrMsgCb(function (\RdKafka\Producer $producer, \RdKafka\Message $message): void {
            if ($message->err !== RD_KAFKA_RESP_ERR_NO_ERROR) {
                throw new \RuntimeException(
                    "Message delivery failed: " . rd_kafka_err2str($message->err)
                );
            }
        });

        $this->producer = new \RdKafka\Producer($config);
        $this->topic = $this->producer->newTopic($topicName);
    }

    /**
     * Send event to Kafka
     *
     * @param string $key Partition key (e.g., user_id)
     * @param array $payload Event data
     */
    public function send(string $key, array $payload): void
    {
        $message = json_encode($payload, JSON_THROW_ON_ERROR);

        $this->topic->producev(
            partition: RD_KAFKA_PARTITION_UA, // Auto-assign partition by key
            msgflags: 0,
            payload: $message,
            key: $key,
            headers: [
                'event_type' => $payload['type'] ?? 'unknown',
                'timestamp' => (string) time(),
                'source' => 'order-service',
            ],
        );

        // Trigger delivery callbacks
        $this->producer->poll(0);
    }

    /**
     * Flush remaining messages before shutdown
     */
    public function flush(int $timeoutMs = 10000): void
    {
        $result = $this->producer->flush($timeoutMs);
        if ($result !== RD_KAFKA_RESP_ERR_NO_ERROR) {
            throw new \RuntimeException('Failed to flush producer');
        }
    }
}

// Usage
$producer = new KafkaProducer('kafka1:9092,kafka2:9092', 'orders');

$producer->send('user_123', [
    'type' => 'order.created',
    'order_id' => 'ord_abc',
    'user_id' => 'user_123',
    'total' => 99.99,
    'items' => [
        ['product_id' => 'prod_1', 'qty' => 2],
    ],
]);

$producer->flush();
### Consumer: чтение сообщений
<?php

declare(strict_types=1);

/**
 * Kafka consumer with manual offset commit
 */
final class KafkaConsumer
{
    private \RdKafka\KafkaConsumer $consumer;

    public function __construct(
        string $brokers,
        string $groupId,
        array $topics,
    ) {
        $config = new \RdKafka\Conf();
        $config->set('metadata.broker.list', $brokers);
        $config->set('group.id', $groupId);
        $config->set('auto.offset.reset', 'earliest'); // Start from beginning if no offset
        $config->set('enable.auto.commit', 'false');    // Manual commit
        $config->set('max.poll.interval.ms', '300000');
        $config->set('session.timeout.ms', '30000');

        // Rebalance callback
        $config->setRebalanceCb(
            function (\RdKafka\KafkaConsumer $consumer, int $err, array $partitions = null): void {
                match ($err) {
                    RD_KAFKA_RESP_ERR__ASSIGN_PARTITIONS => $consumer->assign($partitions),
                    RD_KAFKA_RESP_ERR__REVOKE_PARTITIONS => $consumer->assign(null),
                    default => throw new \RuntimeException(rd_kafka_err2str($err)),
                };
            }
        );

        $this->consumer = new \RdKafka\KafkaConsumer($config);
        $this->consumer->subscribe($topics);
    }

    /**
     * Main consume loop
     */
    public function run(callable $handler): void
    {
        while (true) {
            $message = $this->consumer->consume(1000); // 1 second timeout

            match ($message->err) {
                RD_KAFKA_RESP_ERR_NO_ERROR => $this->processMessage($message, $handler),
                RD_KAFKA_RESP_ERR__PARTITION_EOF => null, // End of partition, wait
                RD_KAFKA_RESP_ERR__TIMED_OUT => null,     // No message, wait
                default => throw new \RuntimeException(
                    "Kafka error: " . rd_kafka_err2str($message->err)
                ),
            };
        }
    }

    private function processMessage(\RdKafka\Message $message, callable $handler): void
    {
        try {
            $payload = json_decode($message->payload, true, 512, JSON_THROW_ON_ERROR);
            $handler($payload, $message->key, $message->headers);

            // Commit offset only after successful processing
            $this->consumer->commit($message);
        } catch (\Throwable $e) {
            // Log error, send to DLQ, or retry
            error_log("Failed to process message: " . $e->getMessage());
        }
    }
}

// Usage
$consumer = new KafkaConsumer(
    brokers: 'kafka1:9092,kafka2:9092',
    groupId: 'order-processor',
    topics: ['orders'],
);

$consumer->run(function (array $payload, string $key, ?array $headers): void {
    match ($payload['type']) {
        'order.created' => handleOrderCreated($payload),
        'order.paid' => handleOrderPaid($payload),
        'order.cancelled' => handleOrderCancelled($payload),
        default => null, // Skip unknown events
    };
});
## Паттерны использования Kafka

Dead Letter Queue (DLQ)

Сообщения, которые не удалось обработать, отправляются в отдельный топик:

<?php

declare(strict_types=1);

final class ResilientConsumer
{
    private const MAX_RETRIES = 3;

    public function __construct(
        private readonly KafkaConsumer $consumer,
        private readonly KafkaProducer $dlqProducer,
    ) {}

    public function processWithRetry(array $payload, string $key): void
    {
        $retries = (int) ($payload['_retry_count'] ?? 0);

        try {
            $this->processOrder($payload);
        } catch (\Throwable $e) {
            if ($retries >= self::MAX_RETRIES) {
                // Send to Dead Letter Queue
                $this->dlqProducer->send($key, [
                    ...$payload,
                    '_dlq_reason' => $e->getMessage(),
                    '_dlq_timestamp' => time(),
                    '_original_topic' => 'orders',
                ]);
                return;
            }

            // Retry with incremented counter
            $this->dlqProducer->send($key, [
                ...$payload,
                '_retry_count' => $retries + 1,
            ]);
        }
    }

    private function processOrder(array $payload): void
    {
        // Business logic
    }
}
### Transactional Outbox

Для гарантии согласованности между БД и Kafka:

<?php

declare(strict_types=1);

/**
 * Outbox pattern: save event to DB, then relay to Kafka
 */
final class OutboxPublisher
{
    public function __construct(
        private readonly \PDO $db,
    ) {}

    /**
     * Save order and event in single transaction
     */
    public function createOrder(array $orderData): string
    {
        $this->db->beginTransaction();

        try {
            $orderId = $this->insertOrder($orderData);

            // Save event to outbox table (same transaction)
            $stmt = $this->db->prepare(<<<SQL
                INSERT INTO outbox_events (id, aggregate_type, aggregate_id, event_type, payload, created_at)
                VALUES (:id, 'order', :aggregate_id, 'order.created', :payload, NOW())
            SQL);

            $stmt->execute([
                'id' => uuid_create(),
                'aggregate_id' => $orderId,
                'payload' => json_encode([
                    'order_id' => $orderId,
                    'user_id' => $orderData['user_id'],
                    'total' => $orderData['total'],
                ]),
            ]);

            $this->db->commit();

            return $orderId;
        } catch (\Throwable $e) {
            $this->db->rollBack();
            throw $e;
        }
    }

    private function insertOrder(array $data): string
    {
        // Insert order into orders table
        return 'ord_' . bin2hex(random_bytes(8));
    }
}
## Мониторинг Kafka

Ключевые метрики

Метрика Описание Порог алерта
Consumer Lag Отставание потребителя > 10000 сообщений
Under-replicated Partitions Недореплицированные партиции > 0
Request Latency Задержка запросов p99 > 100ms
Disk Usage Использование диска > 80%
ISR Shrink Rate Частота выпадения из ISR > 0/min

Итоги

  • Kafka -- распределённый коммит-лог с высокой пропускной способностью и отказоустойчивостью
  • Партиции обеспечивают параллелизм и порядок (внутри одной партиции)
  • Consumer groups позволяют масштабировать потребление
  • php-rdkafka предоставляет полноценный доступ к Kafka из PHP
  • Используйте manual commit, DLQ и outbox pattern для надёжной обработки