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







