0Pricing
PHP Academy · Урок

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 — локальная установка не требуется.

Все уроки этого курса

  1. Зачем нужны асинхронные сообщения
  2. Работа с RabbitMQ в PHP
  3. Apache Kafka с PHP
  4. Создание событийно-ориентированных рабочих процессов
← Назад к PHP Academy