Перейти к основному содержимому

Работа с Kafka

В PHP-сервисах Ensi для Kafka используют расширение phprdkafka (уже есть в базовом docker-образе) и три Laravel-пакета, которые закрывают подключение, чтение и запись. Правила именования топиков — в Kafka Design Guide.

Настройка сервиса

ensi/laravel-phprdkafka

Пакет добавляет менеджер подключений к Kafka по аналогии с менеджером БД. В config/kafka.php описывают параметры соединений, отдельно для consumers и producers; можно завести несколько именованных подключений. Дальше из контейнера или через фасад Kafka получают уже сконфигурированные RdKafka\Producer и RdKafka\KafkaConsumer.

Общие настройки (брокеры, авторизация) выносят в connection. В consumer/producer указывают соединение и добавляют специфичные для клиента параметры.

Пример конфигурации:

# config/kafka.php
return [
'connections' => [
'default' => [
'settings' => [
'metadata.broker.list' => env('KAFKA_BROKER_LIST'),
'security.protocol' => env('KAFKA_SECURITY_PROTOCOL', 'plaintext'),
'sasl.mechanisms' => env('KAFKA_SASL_MECHANISMS'),
'sasl.username' => env('KAFKA_SASL_USERNAME'),
'sasl.password' => env('KAFKA_SASL_PASSWORD'),
'log_level' => env('KAFKA_DEBUG', false) ? (string)LOG_DEBUG : (string)LOG_INFO,
'debug' => env('KAFKA_DEBUG', false) ? 'all' : null,
],
'topics' => [
'foobars' => $contour . '.domain.fact.foobars.1'
]
]
],
'consumers' => [
'default' => [
'connection' => 'default',
'additional-settings' => [
'group.id' => env('KAFKA_CONSUMER_GROUP_ID', env('APP_NAME')),
'enable.auto.commit' => true,
'auto.offset.reset' => 'beginning',
],
],
],
'producers' => [
'default' => [
'connection' => 'default',
'additional-settings' => [
'compression.codec' => env('KAFKA_PRODUCER_COMPRESSION_CODEC', 'snappy'),
],
],
],
];

В настройках подключения можно использовать любые параметры rdkafka. В пакете уже заданы значения по умолчанию:

  • group.id — идентификатор consumer group. Задаётся всегда, даже если пока один инстанс: при масштабировании консьюмеры уже будут в одной группе, останется следить за числом партиций.
  • enable.auto.commit — автоматически обновляет offset через некоторое время после получения сообщения.
  • auto.offset.reset — поведение, когда для группы ещё нет сохранённого offset. beginning означает чтение с начала топика.

В соединении перечисляют все топики приложения. В коде ссылаются на ключ из этого списка, а не на полное имя топика:

'topics' => [
'foobars' => $contour . '.domain.fact.foobars.1'
]

ensi/laravel-phprdkafka-consumer

Пакет добавляет artisan-команду kafka:consume topic и конфигурацию обработчиков сообщений.

# config/kafka-consumer.php
return [
'global_middleware' => [ TraceEventKafkaMiddleware::class ],
'stop_signals' => [SIGTERM, SIGINT],

'processors' => [
[
'topic' => 'foobars', // ключ из kafka.connections.<name>.topics
'consumer' => 'default',
'type' => 'action',
'class' => \App\Domain\Kafka\Actions\Listen\ListenOfferAction::class,
'queue' => false,
'consume_timeout' => 5000,
],
]
];

Для каждого топика указывают класс-обработчик и какой consumer из config/kafka.php использовать. Если для топика нужны другие настройки consumer — заводят отдельное подключение в config/kafka.php и ссылаются на него.

Обработчик — класс с методом execute(RdKafka\Message $message):

class ListenOfferAction
{
public function execute(Message $message)
{
// ...
}
}

Можно задавать middleware по аналогии с HTTP:

class TraceEventKafkaMiddleware {
public function handle(Message $message, Closure $next): mixed
{
// ...
return $next($message);
}
}

ensi/laravel-phprdkafka-producer

Обёртка над RdKafka\Producer для рутинной отправки:

$producer = new HighLevelProducer("my-topic", "my-producer");
$producer->sendOne($payload);

HighLevelProducer берёт продюсер из config/kafka.php, отправляет сообщение и делает flush, чтобы оно ушло из буфера в брокер.

Настройка топиков

Важные параметры задают при создании топика (не в consumer/producer):

  • число партиций — верхняя граница параллелизма consumer group;
  • число реплик — надёжность хранения;
  • retention.bytes / retention.ms — сколько данных и как долго хранить.

Топик создаётся в CI/CD при отгрузке сервиса, который им пользуется (источник или потребитель — не важно). Пайплайн:

  1. запускает php artisan kafka:find-not-created-topics — сверяет топики из config/kafka.php с брокером;
  2. если чего-то нет, читает файл настроек топиков из репозитория деплоя;
  3. создаёт отсутствующие топики.

Пример файла настроек:

topics:
- name: prod.all.fct.my-topic.0
partitions: 2
replicas: 1
config:
- name: retention.ms
value: 60480000 # 7 days
- name: retention.bytes
value: 1073741824 # 1 Gb

Если нужный топик не описан в этом файле, пайплайн завершается с ошибкой.

Для локальной Kafka тот же механизм: в воркспейсе есть файл с топиками и скрипт, который их создаёт.