Apache Kafka с PHP
Передавайте высокопроизводительные события с помощью Kafka
«Apache Kafka с PHP» — бесплатный урок PHP Academy на CoddyKit. Это урок 3 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения PHP Academy, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс PHP Academy содержит 4 уроков всего.
Kafka с PHP
Kafka — не очередь задач, а распределённый журнал фиксации, доступный только для добавления записей. Производители добавляют записи в темы, а потребители читают их начиная с собственного смещения и могут повторно воспроизводить историю. Поэтому Kafka отлично подходит для высокопроизводительных потоков событий, источника событий и передачи одного потока нескольким независимым группам потребителей.
В PHP для работы с Kafka используется ext-rdkafka — привязка к проверенной боевой библиотеке C librdkafka.
Журнал, а не очередь
Как меняется подход при переходе от RabbitMQ к Kafka:
- Сообщения не удаляются после получения; они истекают согласно политике хранения — по времени или размеру.
- Каждый потребитель отслеживает собственное смещение — свою позицию в журнале.
- Тема разделена на разделы; порядок гарантируется только внутри одного раздела.
- Несколько групп потребителей независимо читают одну и ту же тему.
Если Вам нужны повторное воспроизведение, рассылка многим читателям или огромная пропускная способность, выбирайте Kafka. Если нужны маршрутизация каждого сообщения и сроки жизни, выбирайте RabbitMQ.
Установка ext-rdkafka
Сначала установите нативную библиотеку, затем расширение PECL и, при необходимости, оболочку более высокого уровня.
# Debian/Ubuntu
apt-get install -y librdkafka-dev
pecl install rdkafka
echo "extension=rdkafka.so" >> php.ini
# Optional ergonomic wrapper
composer require longlang/phpkafka # or arnaud-lb/php-rdkafka-stubs for IDEСоздание записей
Создайте RdKafka\Producer, получите дескриптор темы и вызовите produce(). Ключ определяет, в какой раздел попадёт запись: один и тот же ключ означает один и тот же раздел, поэтому порядок сохраняется. Всегда вызывайте flush() перед завершением, иначе буферизованные записи будут потеряны.
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$producer = new RdKafka\Producer($conf);
$topic = $producer->newTopic('orders');
// RD_KAFKA_PARTITION_UA = let the partitioner choose by key
$key = 'order-42';
$topic->produce(RD_KAFKA_PARTITION_UA, 0, json_encode(['id' => 42]), $key);
$producer->poll(0);
$result = $producer->flush(10000); // wait up to 10s
if ($result !== RD_KAFKA_RESP_ERR_NO_ERROR) {
throw new RuntimeException('Failed to flush');
}Разделение по ключу
Разделение — основа масштабируемости и упорядочивания в Kafka. Разделитель по умолчанию хеширует ключ записи: partition = hash(key) % numPartitions. Правильный выбор ключа важен:
- Ключ
customerId→ все события клиента сохраняют порядок и остаются в одном разделе. - Нулевой ключ → циклическое распределение между разделами (максимальная пропускная способность, порядок не гарантируется).
Уменьшить число разделов темы невозможно, а добавление разделов изменяет распределение хешей, поэтому заранее задайте число разделов для максимального параллелизма.
Группы потребителей и смещения
Используйте высокоуровневый KafkaConsumer с параметром group.id. Kafka распределяет разделы между участниками группы и выполняет перераспределение, когда участники присоединяются или выходят. Каждый участник читает только назначенные ему разделы, что позволяет бесплатно масштабировать систему горизонтально.
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$conf->set('group.id', 'order-emailers');
$conf->set('auto.offset.reset', 'earliest'); // start of log if no offset
$conf->set('enable.auto.commit', 'false'); // we commit manually
$consumer = new RdKafka\KafkaConsumer($conf);
$consumer->subscribe(['orders']);
while (true) {
$msg = $consumer->consume(5000);
if ($msg->err === RD_KAFKA_RESP_ERR_NO_ERROR) {
handle($msg->payload);
$consumer->commit($msg); // commit offset AFTER work
}
}
function handle(string $p): void {}Когда фиксировать смещение
Фиксация смещения определяет семантику доставки:
- Фиксировать после обработки → доставка как минимум один раз: сбой до фиксации приводит к повторному чтению записи.
- Фиксировать до обработки → доставка не более одного раза: при сбое запись будет потеряна.
Автоматическая фиксация (enable.auto.commit=true) выполняется по таймеру независимо от того, завершилась ли Ваша работа. Это удобно, но при сбое может незаметно привести к потере записей. Отключите её и фиксируйте смещения вручную, когда важна корректность.
Обработка ошибок получения сообщений
Не каждый результат вызова consume() является сообщением. Необходимо проверять код ошибки: __PARTITION_EOF и __TIMED_OUT — штатные управляющие сигналы, а не сбои.
<?php
$msg = $consumer->consume(2000);
switch ($msg->err) {
case RD_KAFKA_RESP_ERR_NO_ERROR:
echo "Got: {$msg->payload} @ offset {$msg->offset}\n";
break;
case RD_KAFKA_RESP_ERR__PARTITION_EOF:
echo "Reached end of partition\n"; // caught up, keep polling
break;
case RD_KAFKA_RESP_ERR__TIMED_OUT:
echo "No message this poll\n";
break;
default:
throw new \Exception($msg->errstr(), $msg->err);
}Настройка пропускной способности
Kafka обеспечивает высокую пропускную способность благодаря пакетной обработке. Основные настройки производителя:
linger.ms— небольшая задержка для объединения большего числа записей в один запрос (например, 5–20 мс).batch.size/queue.buffering.max.messages— более крупные буферы и меньшее число обращений к сети.compression.type—lz4илиzstdзначительно снижают затраты на передачу данных.acks—allдля надёжности,1для меньшей задержки.
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$conf->set('compression.type', 'lz4');
$conf->set('linger.ms', '10');
$conf->set('batch.size', '65536');
$conf->set('acks', 'all');
$producer = new RdKafka\Producer($conf);Обратные вызовы отчёта о доставке
Поскольку produce() работает асинхронно, неудачная отправка не вызовет исключение сразу. Зарегистрируйте обратный вызов отчёта о доставке в конфигурации производителя, чтобы узнать судьбу каждой записи. В PHP это единственный надёжный способ обнаружить незаметные сбои производителя.
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$conf->setDrMsgCb(function ($producer, $msg) {
if ($msg->err) {
fwrite(STDERR, 'Delivery FAILED: ' . rd_kafka_err2str($msg->err) . "\n");
} else {
echo "Delivered to partition {$msg->partition} @ offset {$msg->offset}\n";
}
});
$producer = new RdKafka\Producer($conf);
// poll() services the callback queue; call it after producing
$producer->poll(0);Особенности PHP
Kafka рассчитана на долгоживущие клиенты, а жизненный цикл запроса PHP этому препятствует:
- Производители асинхронно накапливают данные в буфере — всегда вызывайте
flush()перед завершением скрипта, иначе записи будут потеряны. - Запускайте потребителей как постоянные рабочие процессы CLI под управлением супервизора, а не внутри веб-запроса.
- Перераспределение приостанавливает получение сообщений; сокращайте обработку каждого сообщения или задавайте для
max.poll.interval.msдостаточно большое значение, чтобы потребителя не исключили из группы. - Задайте
log_levelи зарегистрируйте обратный вызов отчёта о доставке производителя, чтобы обнаруживать незаметные сбои.
Быстрая проверка
Гарантии порядка в Kafka.
Повторение
Kafka в PHP:
- Kafka — это журнал, который можно воспроизводить, а не очередь; потребители отслеживают смещения.
- Разделение по ключу обеспечивает порядок для каждого ключа и масштабирование.
- Группы потребителей делят разделы между собой и автоматически выполняют перераспределение.
- Фиксируйте смещения после работы, чтобы обеспечить доставку как минимум один раз; отключайте автоматическую фиксацию для полного контроля.
- Настраивайте
linger.ms,batch.sizeиcompression.type; всегда вызывайтеflush().
Далее: объединение этих примитивов в надёжные рабочие процессы на основе событий.
Часто задаваемые вопросы
Урок «Apache Kafka с PHP» бесплатный?
Да — полный текст урока «Apache Kafka с PHP» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс PHP Academy, подпишись на CoddyKit PRO. Курс PHP Academy содержит 4 уроков всего.
Чему я научусь в уроке «Apache Kafka с PHP»?
Передавайте высокопроизводительные события с помощью Kafka Ты практикуешь PHP Academy с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать PHP Academy?
Предыдущий опыт не требуется. PHP Academy на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 3 из 4.
Сколько времени занимает урок «Apache Kafka с PHP»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке PHP Academy?
Да. Каждый урок PHP Academy включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Зачем нужны асинхронные сообщения
- Работа с RabbitMQ в PHP
- Apache Kafka с PHP
- Создание событийно-ориентированных рабочих процессов