Очереди сообщений — критический компонент в микросервисной архитектуре. Apache Kafka (Wikipedia) справляется с нагрузками от 100 тысяч сообщений в секунду, но неправильная настройка приводит к потере данных или недетерминированному порядку. Например, в одном проекте клиент потерял заказы из-за установки acks=0: продюсер не дожидался подтверждения, и при сбое брокера сообщения исчезали бесследно. После миграции на acks=all с enable.idempotence=true проблема исчезла — доставка стала гарантированной, а порядок событий по одному ключу сохранился. Кроме того, мы внедрили Schema Registry для контроля версий сообщений: теперь разные сервисы могут безопасно эволюционировать свои данные. Настроим Kafka под ваш проект за 3–5 дней с полной документацией.
Почему Kafka, а не RabbitMQ для потоковой обработки?
Kafka хранит сообщения по retention-политике (например, 7 дней или 100 ГБ), позволяя множеству consumer'ов читать один топик с разными смещениями. RabbitMQ удаляет сообщения после подтверждения — это хорошо для task-очередей, но плохо для аудита и replay. В потоковых сценариях Kafka примерно в 10 раз производительнее RabbitMQ по пропускной способности. При этом Kafka в 3 раза экономичнее на хранении большого объёма данных за счёт последовательной записи и сжатия.
| Критерий | Apache Kafka | RabbitMQ |
|---|---|---|
| Хранение сообщений | По retention (фиксированное время/размер) | До подтверждения (ack) |
| Повторное чтение (replay) | Поддерживается (сброс offset) | Нет |
| Потоковая обработка | Встроенная (Kafka Streams) | Требует внешних инструментов |
| Максимальная пропускная способность | Миллионы сообщений/с | Сотни тысяч сообщений/с |
| Типичное применение | Event sourcing, аналитика, логи | Task очереди, RPC, уведомления |
Пропускная способность Kafka достигает 2–3 миллионов сообщений в секунду на кластере из трёх нод, тогда как RabbitMQ — не более 300–500 тысяч. Retention по умолчанию — 7 дней, но мы настраиваем под бизнес-логику, например, для аудита — 30 дней, для логов — 100 ГБ.
Как мы настраиваем Kafka: стек, конфиги, процесс
Используем Confluent-образы Kafka 7.6+ с обязательным Schema Registry для контракта данных. Для PHP применяем librdkafka, для Node.js — kafkajs. Процесс:
- Анализ нагрузки и проектирование топиков (количество партиций, фактор репликации).
- Развёртывание через Docker Compose или Kubernetes (Strimzi).
- Настройка producer с idempotent и acks=all.
- Реализация consumer'ов с ручным commit after processing.
- Мониторинг через JMX + Grafana с алертами при lag > 10 000.
Установка через Docker
# docker-compose.yml services: zookeeper: image: confluentinc/cp-zookeeper:7.6.0 environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 volumes: - zookeeper_data:/var/lib/zookeeper/data - zookeeper_log:/var/lib/zookeeper/log kafka: image: confluentinc/cp-kafka:7.6.0 depends_on: [zookeeper] environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false" KAFKA_LOG_RETENTION_HOURS: 168 # 7 дней KAFKA_LOG_RETENTION_BYTES: 107374182400 # 100 GB KAFKA_NUM_PARTITIONS: 6 KAFKA_DEFAULT_REPLICATION_FACTOR: 1 volumes: - kafka_data:/var/lib/kafka/data ports: - "9092:9092" kafka-ui: image: provectuslabs/kafka-ui:latest depends_on: [kafka] environment: KAFKA_CLUSTERS_0_NAME: local KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:9092 ports: - "8080:8080" volumes: zookeeper_data: zookeeper_log: kafka_data: Создание топиков
# Создать топик с 6 партициями и репликой 1 (для одиночного брокера) kafka-topics.sh --bootstrap-server kafka:9092 \ --create \ --topic user-events \ --partitions 6 \ --replication-factor 1 \ --config retention.ms=604800000 \ --config cleanup.policy=delete # Просмотр kafka-topics.sh --bootstrap-server kafka:9092 --describe --topic user-events PHP: Producer (librdkafka)
use RdKafka\Producer; use RdKafka\Conf; class KafkaProducer { private Producer $producer; public function __construct() { $conf = new Conf(); $conf->set('bootstrap.servers', config('kafka.brokers')); $conf->set('security.protocol', 'PLAINTEXT'); $conf->set('acks', 'all'); // подтверждение от всех реплик $conf->set('retries', '3'); $conf->set('enable.idempotence', 'true'); // ровно одна запись $conf->set('compression.type', 'snappy'); $conf->setDrMsgCb(function ($kafka, $message) { if ($message->err !== RD_KAFKA_RESP_ERR_NO_ERROR) { Log::error('Kafka delivery failed', [ 'error' => $message->errstr(), 'topic' => $message->topic_name, ]); } }); $this->producer = new Producer($conf); } public function publish(string $topic, string $key, array $payload): void { $rdTopic = $this->producer->newTopic($topic); $rdTopic->produce( partition: RD_KAFKA_PARTITION_UA, // автовыбор партиции по key msgflags: 0, payload: json_encode($payload), key: $key, // один ключ → одна партиция → порядок событий ); $this->producer->poll(0); } public function flush(): void { $result = $this->producer->flush(10000); // 10 секунд таймаут if (RD_KAFKA_RESP_ERR_NO_ERROR !== $result) { throw new \RuntimeException('Kafka flush failed: ' . rd_kafka_err2str($result)); } } } // Использование $producer->publish('user-events', (string) $user->id, [ 'event' => 'user.registered', 'user_id' => $user->id, 'email' => $user->email, 'timestamp' => now()->toIso8601String(), ]); $producer->flush(); Node.js: Producer и Consumer на kafkajs
import { Kafka, CompressionTypes } from 'kafkajs'; const kafka = new Kafka({ clientId: 'myapp-api', brokers: [process.env.KAFKA_BROKERS!], retry: { retries: 5, initialRetryTime: 300, factor: 0.2, }, }); // Producer const producer = kafka.producer({ allowAutoTopicCreation: false, idempotent: true, maxInFlightRequests: 5, }); await producer.connect(); await producer.send({ topic: 'user-events', compression: CompressionTypes.Snappy, messages: [{ key: String(userId), value: JSON.stringify({ event: 'user.login', userId, ip, timestamp: Date.now() }), headers: { 'content-type': 'application/json' }, }], }); // Consumer const consumer = kafka.consumer({ groupId: 'audit-service' }); await consumer.connect(); await consumer.subscribe({ topic: 'user-events', fromBeginning: false }); await consumer.run({ eachMessage: async ({ topic, partition, message }) => { const payload = JSON.parse(message.value!.toString()); await AuditLog.create({ event: payload.event, userId: payload.userId, metadata: payload, }); }, }); Как добиться идемпотентности producer?
Идемпотентность означает, что запись сообщения в топик происходит ровно один раз, даже при повторной отправке. Для этого в Kafka нужно включить enable.idempotence=true и acks=all. При этом продюсер получает уникальный идентификатор (producer id и sequence number), а брокер отбрасывает дубликаты. Это обязательная настройка для финансовых операций и аудита.
Как настроить мониторинг лага consumer?
Для мониторинга используйте команду kafka-consumer-groups.sh --describe --bootstrap-server localhost:9092 --group <group_name>. Она показывает смещение consumer'а (current-offset) и последнее сообщение в топике (log-end-offset). Разница — это lag. Для автоматизации мы устанавливаем JMX Exporter, который публикует метрики в Prometheus, а затем строим дашборды в Grafana с порогом алерта при lag > 10 000. Это позволяет вовремя заметить отставание и перераспределить партиции.
Что входит в настройку под ключ
- Аудит текущей архитектуры — анализ нагрузки, выбор количества топиков и партиций.
- Развёртывание кластера — Docker Compose или Kubernetes (Strimzi) с учётом отказоустойчивости.
- Интеграция с приложением — написание producer/consumer на PHP или Node.js с обработкой ошибок.
- Schema Registry — внедрение Avro-схем для контроля версий сообщений.
- Мониторинг — дашборд Grafana с алертами по лагу и загрузке брокеров.
- Обучение команды — документация по эксплуатации и аварийному восстановлению.
Сроки реализации
| Этап | Срок |
|---|---|
| Базовый кластер + producer/consumer (один язык) | 3–4 дня |
| Schema Registry + Avro | +2 дня |
| Kafka Streams для агрегации | 3–5 дней |
| Кластер на 3 ноды в Kubernetes | 4–5 дней |
Стоимость рассчитывается индивидуально — зависит от сложности интеграции и необходимости Streams/KSQL. Мы работаем с Kafka более 5 лет и реализовали 20+ проектов по настройке очередей. Получите консультацию по настройке Kafka — оценим задачу за 1 день.







