Обёртка над ext-rdkafka для работы с Apache Kafka на PHP. Библиотека предоставляет простой и удобный интерфейс для продюсирования и консьюминга сообщений.
- PHP 8.4+
- ext-rdkafka
composer require anktx/kafka-clientuse 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();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);
}
}Для более чистого кода используйте генератор:
use Anktx\Kafka\Client\KafkaMessageStream;
$stream = new KafkaMessageStream($consumer);
foreach ($stream->stream() as $message) {
// Только сообщения, без обработки таймаутов/EOF
echo $message->body . "\n";
$consumer->commit($message);
}По умолчанию поток переживает недоступность брокеров (librdkafka
переподключается в фоне): из consume() любой обрыв виден только как серия
таймаутов. Реакция на нештатные ситуации — инжектируемый наблюдатель
StreamObserver: каждый результат consume() (сообщение, таймаут, EOF)
передаётся его хукам onMessage/onTimeout/onEof до выдачи сообщения
наружу, исключение из хука прерывает генератор. Дефолт
SilentStreamObserver поглощает всё — прежнее поведение.
Типовой fail-fast watchdog — бюджет тишины: если сообщения непрерывно
отсутствуют дольше порога (наращивается в onTimeout, сбрасывается в
onMessage), генератор прерывается исключением из хука, воркер падает,
супервизор пересоздаёт процесс (restart-политика Docker, restartPolicy
Kubernetes). Порог — продуктовое решение приложения: потеря брокеров из
consume() неотличима от тишины в топике, поэтому готового класса в
библиотеке нет. Свои сценарии — реализуйте интерфейс StreamObserver
или наследуйте SilentStreamObserver, переопределяя только нужные хуки.
При отправке сообщений они попадают в локальную очередь, а затем асинхронно отправляются в 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-20Psr\Clock\ClockInterface, по умолчанию системные часы)ProbabilityPollStrategy— вызыватьpoll()с вероятностью N (например, 10% вызовов)
Ошибки доставки сообщений (delivery reports) продюсер логирует через PSR-3:
успешная доставка — debug, превышение message.timeout.ms — warning
(ожидаемое следствие недоступности брокеров), прерывание доставки при
остановке продюсера — info, прочие сбои — error. Отчёты доставляются
callback'ам только при вызове poll(): со стратегиями опроса — в фоне,
с NeverPollStrategy — только в момент flush().
$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
);$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.
Дефолты библиотеки ужесточены относительно 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 при недоставленной очереди).
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 мин) для продакшена непригодны.
maxPollIntervalMs — максимальная задержка между вызовами consume().
Если обработка пакета сообщений занимает дольше — консьюмер покидает группу
(партиции отзываются, rebalance), а следующий commit упирается в fencing.
Контракт потребителя библиотеки: время обработки батча между poll'ами +
коммит ≤ maxPollIntervalMs. Тяжёлая обработка — выносите наружу
(бустер-очередь, acknowledge после завершения), не увеличивая таймаут
без нужды: он же работает как watchdog зависшего воркера.
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— таймаут (нет новых сообщений)KafkaPartitionEof— достигнут конец партиции
Недоступность брокеров отдельного результата не образует: событие
ALL_BROKERS_DOWN перехватывается librdkafka и уходит в error-callback,
поэтому из consume() любой обрыв выглядит как серия таймаутов (быстрый
сигнал активного обрыва — warning в error-callback,
RD_KAFKA_RESP_ERR__TRANSPORT; «молчаливый» обрыв — только в логах
librdkafka через ~2 минуты). Прочие коды ошибок RD_KAFKA_RESP_ERR__*
(фатальные, превышение max.poll.interval.ms, fetch-ошибки) бросают
KafkaConsumerException.
Пример обработки — match по классу даёт исчерпывающую диспетчеризацию
(при появлении нового варианта union будет UnhandledMatchError, а не
молчаливый пропуск) и сужение типа внутри веток:
$result = $consumer->consume(1000);
match ($result::class) {
KafkaConsumerMessage::class => $consumer->commit($result),
KafkaConsumeTimeout::class => null, // нет сообщений, можно продолжить работу
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/ # Результаты консьюминга
│ ├── 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/onEof)
│ └── SilentStreamObserver.php # Молчаливая реакция (по умолчанию, расширяемая база)
│
├── Topic/ # Топики
│ ├── Topic.php # Имя топика (VO: непустая строка)
│ └── TopicList.php # Список топиков для подписки
│
├── KafkaConsumer.php # Главный класс консьюмера
├── KafkaProducer.php # Главный класс продюсера
└── KafkaMessageStream.php # Генератор для стриминга сообщений
Библиотека использует иерархию из двух семейств и маркерного интерфейса:
KafkaClientException # Маркер: всё, что кидает библиотека (interface, extends \Throwable)
├── KafkaException # Сбои Kafka/окружения (наследует RdKafka\Exception)
│ ├── 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.