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







