anktx/kafka-client

PHP wrapper for RdKafka

Maintainers

Package info

github.com/anktx/kafka-client

pkg:composer/anktx/kafka-client

Transparency log

Statistics

Installs: 115

Dependents: 0

Suggesters: 0

Stars: 0

Open Issues: 0

0.13.0 2026-08-24 12:17 UTC

This package is auto-updated.

Last update: 2026-08-24 12:24:02 UTC


README

Обёртка над ext-rdkafka для работы с Apache Kafka на PHP. Библиотека предоставляет простой и удобный интерфейс для продюсирования и консьюминга сообщений.

Требования

  • PHP 8.4+
  • ext-rdkafka

Установка

composer require anktx/kafka-client

Быстрый старт

Producer

use Anktx\Kafka\Client\Config\Brokers;
use Anktx\Kafka\Client\Config\ProducerConfig;
use Anktx\Kafka\Client\Config\Enum\CompressionType;
use Anktx\Kafka\Client\KafkaProducer;
use Anktx\Kafka\Client\KafkaMessage\KafkaProducerMessage;
use Anktx\Kafka\Client\Topic\Topic;

$producer = new KafkaProducer(
    new ProducerConfig(
        brokers: new Brokers('kafka:9092'),
        compressionType: CompressionType::Snappy,
    )
);

$producer->produce(
    new KafkaProducerMessage(
        topic: new Topic('events'),
        body: json_encode(['event' => 'order_created', 'id' => 123]),
        key: 'order-123',
        headers: ['source' => 'api'],
    )
);

$producer->flush();

Consumer

use Anktx\Kafka\Client\Config\Brokers;
use Anktx\Kafka\Client\Config\ConsumerConfig;
use Anktx\Kafka\Client\Config\Enum\OffsetReset;
use Anktx\Kafka\Client\KafkaConsumer;
use Anktx\Kafka\Client\KafkaMessage\KafkaConsumerMessage;
use Anktx\Kafka\Client\Topic\Topic;
use Anktx\Kafka\Client\Topic\TopicList;

$consumer = new KafkaConsumer(
    new ConsumerConfig(
        brokers: new Brokers('kafka:9092'),
        groupId: 'order-processor',
        instanceId: 'worker-1',
        offsetReset: OffsetReset::Latest,
    )
);

$consumer->subscribe(
    TopicList::create(new Topic('events'))
);

while (true) {
    $result = $consumer->consume();

    if ($result instanceof KafkaConsumerMessage) {
        echo $result->body . "\n";
        // ... обработка сообщения ...

        $consumer->commit($result);
    }
}

Message Stream

Для более чистого кода используйте генератор:

use Anktx\Kafka\Client\KafkaMessageStream;

$stream = new KafkaMessageStream($consumer);

foreach ($stream->stream() as $message) {
    // Только сообщения, без обработки таймаутов/EOF
    echo $message->body . "\n";
    $consumer->commit($message);
}

По умолчанию поток переживает полную потерю брокеров (librdkafka переподключается в фоне). Реакция на нештатные ситуации — инжектируемый наблюдатель StreamObserver: каждый результат consume() (сообщение, таймаут, потеря всех брокеров, EOF) передаётся его хукам onMessage/onTimeout/onBrokersDown/onEof до выдачи сообщения наружу, исключение из хука прерывает генератор. Дефолт SilentStreamObserver поглощает всё — прежнее поведение.

Готовая fail-fast реализация — BrokersDownBudgetStreamObserver: если брокеры недоступны непрерывно дольше maxBrokersDownMs (wall-clock, источник времени — PSR-20 Psr\Clock\ClockInterface, по умолчанию системные часы), генератор выбрасывает KafkaBrokersDownException — воркер падает, супервизор пересоздаёт процесс (restart-политика Docker, restartPolicy Kubernetes):

use Anktx\Kafka\Client\StreamObserver\BrokersDownBudgetStreamObserver;

$stream = new KafkaMessageStream(
    $consumer,
    new BrokersDownBudgetStreamObserver(maxBrokersDownMs: 30_000),
);

Сообщение и EOF доказывают живое соединение и сбрасывают бюджет; таймаут — нет (не отличает тишину в топике от сетевой проблемы). Свои сценарии реакции — реализуйте интерфейс StreamObserver.

Стратегии опроса (Poll Strategies)

При отправке сообщений они попадают в локальную очередь, а затем асинхронно отправляются в Kafka. Метод poll() обслуживает эту очередь — обрабатывает отчёты о доставке и освобождает память. Если не вызывать poll(), очередь может переполниться.

Стратегии определяют, когда вызывать poll():

use Anktx\Kafka\Client\PollStrategy\TimeoutPollStrategy;
use Anktx\Kafka\Client\PollStrategy\ProbabilityPollStrategy;

// Опрос не чаще, чем раз в N миллисекунд
$producer = new KafkaProducer(
    $config,
    new TimeoutPollStrategy(pollIntervalMs: 1000),
);

// Опрос с вероятностью N (0.0 - 1.0)
$producer = new KafkaProducer(
    $config,
    new ProbabilityPollStrategy(probability: 0.1),
);

Доступные стратегии:

  • NeverPollStrategy — не вызывать poll() (по умолчанию, подходит для низкой нагрузки)
  • TimeoutPollStrategy — вызывать poll() с фиксированным интервалом в миллисекундах (источник времени — PSR-20 Psr\Clock\ClockInterface, по умолчанию системные часы)
  • ProbabilityPollStrategy — вызывать poll() с вероятностью N (например, 10% вызовов)

Ошибки доставки сообщений (delivery reports) продюсер логирует через PSR-3 на уровне error. Отчёты доставляются callback'ам только при вызове poll(): со стратегиями опроса — в фоне, с NeverPollStrategy — только в момент flush().

Конфигурация

ProducerConfig

$config = new ProducerConfig(
    brokers: Brokers,                   // Обязательно (VO: host[:port][,...] )
    queueBufferingMaxKBytes: int,       // По умолчанию: 20480
    batchSize: int,                     // По умолчанию: 102400
    lingerMs: int,                      // По умолчанию: 10
    compressionType: CompressionType,   // По умолчанию: snappy; none — отключить сжатие
    enableIdempotence: bool,            // По умолчанию: true
    messageTimeoutMs: int,              // По умолчанию: 120000
    connectionsMaxIdleMs: int,          // По умолчанию: 180000
    reconnectBackoffMs: int,            // По умолчанию: 100
    reconnectBackoffMaxMs: int,         // По умолчанию: 10000
    socketKeepaliveEnable: bool,        // По умолчанию: true
    isDebug: bool,                      // По умолчанию: false
);

ConsumerConfig

$config = new ConsumerConfig(
    brokers: Brokers,                   // Обязательно (VO: host[:port][,...])
    groupId: string,                    // Обязательно
    instanceId: ?string,                // По умолчанию: null
    offsetReset: OffsetReset,           // По умолчанию: earliest
    autoCommitMs: ?int,                 // По умолчанию: null (ручной коммит)
    sessionTimeoutMs: int,              // По умолчанию: 30000
    heartbeatIntervalMs: int,           // По умолчанию: 3000
    maxPollIntervalMs: int,             // По умолчанию: 300000
    connectionsMaxIdleMs: int,          // По умолчанию: 540000
    reconnectBackoffMs: int,            // По умолчанию: 100
    reconnectBackoffMaxMs: int,         // По умолчанию: 10000
    socketKeepaliveEnable: bool,        // По умолчанию: true
    isDebug: bool,                      // По умолчанию: false
);

Невалидные значения отбрасываются в конструкторах Brokers (пустой список или запись вне формата host[:port]), конфигов (пустой groupId, неположительные интервалы, reconnectBackoffMaxMs < reconnectBackoffMs, heartbeatIntervalMs > sessionTimeoutMs / 3) и сообщений/подписок (пустое имя топика — Topic) исключением InvalidConfigException / InvalidTopicException.

Production defaults: сокеты и надёжность

Дефолты библиотеки ужесточены относительно librdkafka под продовый сценарий «молчаливого» отвала брокера: при обрыве сети без RST (NAT/conntrack выкинул state, firewall дропает пакеты) образуется half-open TCP, и клиент видит лишь RD_KAFKA_RESP_ERR__TIMED_OUT — неотличимо от тишины в топике. У librdkafka из коробки нет ни одного механизма детекта такого соединения: keepalive выключен, connections.max.idle.ms = 0 (не закрывать никогда). У Java-клиента idle-close есть, но занимает 9 минут.

Механизмом детекта в наших дефолтах является TCP keepalive (ядро закрывает мёртвое соединение за ~105 с при рекомендованных ниже sysctl) плюс sessionTimeoutMs для группы и messageTimeoutMs для доставки. connections.max.idle.ms — гигиена соединений, а не детект, поэтому значения для консьюмера и продюсера разные (см. rationale).

Параметр Default Rationale
socketKeepaliveEnable true TCP keepalive на сокетах к брокерам — главный механизм детекта half-open соединения силами ядра; дефолт librdkafka false
connectionsMaxIdleMs (consumer) 540000 Дефолт Java-клиента (9 мин — чуть ниже 10-минутного idle брокера, чтобы клиент закрывал соединение первым). Документация Apache: для classic-консьюмера idle ≥ max.poll.interval.ms (rebalance timeout), иначе соединение с координатором закроется посреди ребаланса большой группы. Соединениям консьюмера детект не нужен: их держит живыми трафик (heartbeat каждые 3 с, постоянные fetch-запросы)
connectionsMaxIdleMs (producer) 180000 Официальная рекомендация Azure Event Hubs (must be < 240000 — их idle-таймаут); укладывается и в 350-секундный idle NLB у AWS MSK. Простаивающий продюсер сам закрывает соединения до того, как это сделает инфраструктурный балансировщик, — не будет записи в уже закрытый сокет
reconnectBackoffMs / MaxMs 100 / 10000 Равны дефолтам librdkafka, зафиксированы явно: экспоненциальный бэкофф реконнекта предсказуем
sessionTimeoutMs 30000 Для static membership (group.instance.id): мёртвый инстанс держит партиции до истечения session timeout; 30 с — баланс скорости failover и ложных срабатываний (рекомендованный Confluent диапазон 6–45 с). Heartbeat шлёт фоновый поток librdkafka, паузы PHP-воркера на сессию не влияют
heartbeatIntervalMs 3000 Дефолт обоих эталонных клиентов (Java и librdkafka): 10 попыток за сессию вместо 3 при session/3 — устойчивость к одиночным задержкам сети и быстрая реакция на ребалансы. Правило документации Kafka «heartbeat ≤ session/3» проверяется в конструкторе (3 × heartbeat ≤ session)
maxPollIntervalMs 300000 Экспозиция контракта «обработка между poll'ами должна укладываться» — см. ниже
enableIdempotence (producer) true Для event transport дубликаты хуже задержек; KIP-679 сделал idempotence дефолтом Java-клиента 3.0+. acks/message.send.max.retries вручную не выставляйте — librdkafka подберёт совместимые значения сам
messageTimeoutMs (producer) 120000 Соответствует delivery.timeout.ms Java-клиента (его самый обкатанный пресет): зависшая доставка вскрывается delivery report'ом за 2 минуты вместо 5

messageTimeoutMs задаёт бюджет доставки одного сообщения силами librdkafka. Прикладной таймаут flush() у потребителя библиотеки должен быть не меньше messageTimeoutMs — иначе flush закончится раньше, чем librdkafka перестанет ретраить, и статус доставки станет неизвестен (KafkaFlushTimeoutException при недоставленной очереди).

Требования к хосту: TCP keepalive

socket.keepalive.enable включает SO_KEEPALIVE, но тайминги проб задаются системными sysctl на хосте (в контейнере — на воркер-ноде). Рекомендуемые значения:

net.ipv4.tcp_keepalive_time=60
net.ipv4.tcp_keepalive_intvl=15
net.ipv4.tcp_keepalive_probes=3

С такими таймингами мёртвое соединение ядро закрывает за ~60+3×15=105 с — детект срабатывает, даже когда приложение полностью молчит: keepalive-таймер отсчитывается от последнего пакета в соединении. Если же приложение пишет в мёртвое соединение, детект берут на себя TCP-ретрансмиссии (tcp_retries2 ≈ 15–30 мин). Побочный эффект: пробы сами являются трафиком не реже раза в 60 с — живое простаивающее соединение не выглядит простаивающим для idle-таймаутов балансировщиков. Дефолты Linux (7200/75/9 = 2 ч 11 мин) для продакшена непригодны.

Контракт max.poll.interval.ms

maxPollIntervalMs — максимальная задержка между вызовами consume(). Если обработка пакета сообщений занимает дольше — консьюмер покидает группу (партиции отзываются, rebalance), а следующий commit упирается в fencing. Контракт потребителя библиотеки: время обработки батча между poll'ами + коммит ≤ maxPollIntervalMs. Тяжёлая обработка — выносите наружу (бустер-очередь, acknowledge после завершения), не увеличивая таймаут без нужды: он же работает как watchdog зависшего воркера.

OffsetReset: политика при отсутствии закоммиченного смещения

offsetReset отвечает на вопрос «с чего начать чтение, если у группы нет валидного закоммиченного смещения». Такое бывает, когда:

  • группа новая и коммитов ещё не было (частный случай — опечатка в groupId);
  • закоммиченный офсет устарел: данные под ним удалены retention-политикой или истёк срок хранения офсетов группы (offsets.retention.minutes).

Если валидный офсет есть — политика не активируется вовсе, все три значения работают одинаково.

Кейс Значение для librdkafka Поведение
OffsetReset::Earliest earliest сброс на начало партиции (перечитает всю историю)
OffsetReset::Latest latest сброс на конец (молча пропустит всю историю)
OffsetReset::Error error сброс запрещён: партиция переходит в ошибку RD_KAFKA_RESP_ERR__AUTO_OFFSET_RESET, consume() бросает KafkaConsumerException

OffsetReset::Error — это strict-режим: опечатка в groupId, потерянная история или невалидный офсет не проходят молча, а останавливают цикл потребления исключением.

Почему Error, а не None. Одна и та же политика в разных клиентах называется по-разному: Java-клиент и документация Kafka называют её none, а librdkafka (и ext-rdkafka) — error; значение none librdkafka отвергает как невалидное (Invalid value "none" for configuration property "auto.offset.reset"). Библиотека работает через librdkafka, поэтому кейс назван по его канону: Error = 'error'. Отображение имён — в таблице выше.

Логирование

PSR-3 логгер передаётся напрямую в конструкторы клиентов (по умолчанию NullLogger):

use Psr\Log\LoggerInterface;

$producer = new KafkaProducer($config, logger: $logger);
$consumer = new KafkaConsumer($config, logger: $logger);

Типы возвращаемых значений

Метод consume() возвращает union type (все варианты реализуют интерфейс ConsumeResult):

  • KafkaConsumerMessage — успешно полученное сообщение
  • KafkaConsumeTimeout — таймаут (нет новых сообщений)
  • KafkaBrokersDown — полная потеря соединения со всеми брокерами (ALL_BROKERS_DOWN): не ошибка — librdkafka переподключается в фоновых потоках; отдельный результат для метрик и watchdog'ов
  • KafkaPartitionEof — достигнут конец партиции

Пример обработки — match по классу даёт исчерпывающую диспетчеризацию (при появлении нового варианта union будет UnhandledMatchError, а не молчаливый пропуск) и сужение типа внутри веток:

$result = $consumer->consume(1000);

match ($result::class) {
    KafkaConsumerMessage::class => $consumer->commit($result),
    KafkaConsumeTimeout::class => null, // нет сообщений, можно продолжить работу
    KafkaBrokersDown::class => null,    // все брокеры недоступны, librdkafka переподключается
    KafkaPartitionEof::class => null,   // достигнут конец партиции
};

ConsumeResult — именованный supertype для хелперов, логгеров и метрик, принимающих любой результат consume() целиком. Если нужна только обработка сообщений без ветвления — см. Message Stream.

Структура проекта

src/
├── Clock/                           # Время
│   └── SystemClock.php              # Системные часы PSR-20 (по умолчанию)
│
├── Config/                          # Конфигурация
│   ├── Brokers.php                  # Список брокеров (VO: host[:port][,...] )
│   ├── ConsumerConfig.php           # Конфигурация консьюмера
│   ├── ProducerConfig.php           # Конфигурация продюсера
│   └── Enum/                        # Перечисления
│       ├── CompressionType.php      # Типы компрессии (none, snappy, gzip, lz4, zstd)
│       └── OffsetReset.php          # Стратегия сброса оффсета (earliest, latest, error)
│
├── ConsumeResult/                   # Результаты консьюминга
│   ├── KafkaBrokersDown.php        # Полная потеря всех брокеров (ALL_BROKERS_DOWN)
│   ├── KafkaConsumeTimeout.php      # Таймаут (нет сообщений)
│   └── KafkaPartitionEof.php        # Достигнут конец партиции
│
├── Exception/                       # Исключения
│   ├── KafkaClientException.php     # Маркерный интерфейс всех исключений библиотеки
│   ├── Kafka/                       # Сбои инфраструктуры Kafka (runtime)
│   └── Logic/                       # Детерминированные ошибки программиста
│
├── KafkaMessage/                    # Сообщения
│   ├── KafkaConsumerMessage.php     # Сообщение консьюмера (topic/partition/offset обязательны)
│   └── KafkaProducerMessage.php     # Сообщение продюсера (валидация в конструкторе)
│
├── Log/                             # Логирование
│   ├── RdKafkaCallbacks.php         # Колбэки librdkafka + единая политика логирования в PSR-3
│   └── RdKafkaLogLevel.php          # Маппинг syslog-severity librdkafka в уровни PSR-3
│
├── PollStrategy/                    # Стратегии опроса очереди
│   ├── PollStrategy.php             # Интерфейс стратегии
│   ├── NeverPollStrategy.php        # Не вызывать poll()
│   ├── ProbabilityPollStrategy.php  # Вызывать с вероятностью N
│   └── TimeoutPollStrategy.php      # Вызывать с фиксированным интервалом (мс)
│
├── StreamObserver/                  # Реакция на результаты consume() в потоке сообщений
│   ├── StreamObserver.php           # Интерфейс наблюдателя (onMessage/onTimeout/onBrokersDown/onEof)
│   ├── SilentStreamObserver.php     # Молчаливая реакция (по умолчанию)
│   └── BrokersDownBudgetStreamObserver.php # Fail-fast: брокеры недоступны дольше maxBrokersDownMs
│
├── Topic/                           # Топики
│   ├── Topic.php                    # Имя топика (VO: непустая строка)
│   └── TopicList.php                # Список топиков для подписки
│
├── KafkaConsumer.php                # Главный класс консьюмера
├── KafkaProducer.php                # Главный класс продюсера
└── KafkaMessageStream.php           # Генератор для стриминга сообщений

Обработка исключений

Библиотека использует иерархию из двух семейств и маркерного интерфейса:

KafkaClientException                 # Маркер: всё, что кидает библиотека (interface, extends \Throwable)
├── KafkaException                   # Сбои Kafka/окружения (наследует RdKafka\Exception)
│   ├── KafkaBrokersDownException    # Брокеры недоступны дольше maxBrokersDownMs (BrokersDownBudgetStreamObserver)
│   ├── KafkaConsumerException        # Ошибка консьюмера
│   ├── KafkaFlushTimeoutException    # Таймаут flush: очередь не отправлена за $timeoutMs
│   └── KafkaProducerException        # Ошибка продюсера
└── LogicException                   # Ошибки программиста (наследует \LogicException)
    ├── ClientClosedException         # Операция после close() клиента
    ├── EmptySubscriptionsException   # Пустой список подписок
    ├── InvalidConfigException        # Невалидная конфигурация или параметры
    ├── InvalidMessageException       # Невалидные свойства сообщения (partition, timestampMs)
    ├── InvalidTopicException         # Пустое имя топика (Topic VO)
    └── NotSubscribedException        # Не подписан на топики

Точки поимки:

  • catch (KafkaClientException) — всё, что кидает библиотека;
  • catch (KafkaException) / catch (RdKafka\Exception) — только сбои Kafka, без опечаток в конфиге;
  • catch (\LogicException) — детерминированные ошибки использования (невалидный конфиг, неверный порядок вызовов).

Пример обработки:

try {
    $producer->produce($message);
    $producer->flush();
} catch (KafkaFlushTimeoutException $e) {
    // Flush не успел за $timeoutMs: часть сообщений могла остаться
    // в локальной очереди продюсера — статус доставки неизвестен
} catch (KafkaProducerException $e) {
    // Ошибка отправки сообщения
} catch (KafkaClientException $e) {
    // Всё остальное от библиотеки
}

Жизненный цикл консьюмера

KafkaConsumer::close() идемпотентен; операции после закрытия бросают ClientClosedException (ошибка программирования, не ретраить). Подробно, с обоснованием и примером шаблона: docs/lifecycle.md.