Очереди сообщений — критический компонент в микросервисной архитектуре. 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 день.







