Kafka-клієнти в продакшені — це не просто «відправити повідомлення» і «отримати повідомлення». Часто від неправильної конфігурації страждає consistency бізнес-даних. Нещодавно ми зіткнулися з проєктом, де через auto-commit та відсутність ідемпотентності втрачалося до 10% замовлень під час перезапуску консьюмера. Після впровадження ручного керування офсетами та ідемпотентного продюсера втрати звели до нуля. Вартість таких втрат для середнього інтернет-магазину може сягати 2 млн грн на місяць. Ми спеціалізуємося на розробці надійних продюсерів та консьюмерів для високонавантажених веб-додатків на Java та Python. У статті розберемо ключові проблеми, типові конфіги та практичні кейси: від налаштування ідемпотентності до паралельної обробки зі збереженням порядку.
Проблеми, які ми вирішуємо
Втрата повідомлень при збоях брокера або консьюмера
Без правильного налаштування acks та ідемпотентності ви ризикуєте втратити критичні для бізнесу дані. Наприклад, при acks=1 повідомлення може бути втрачене, якщо лідер падає до реплікації. Ми використовуємо acks=all з min.insync.replicas=2 та enable.idempotence=true. Це виключає втрати навіть при відмові одного брокера. Економія на інфраструктурі за рахунок відмови від дорогих рішень — до 40%.
Порушення порядку обробки та дублювання
Консьюмери без керування offset'ами можуть обробити одне повідомлення двічі, а паралельна обробка руйнує порядок. Рішення — ручний commit після обробки та роутинг за ключем у чергу всередині партиції. В одному проєкті це скоротило кількість помилок валідації на 30%.
Довгий ребаланс та повторна обробка
При збої консьюмера та повторному ребалансі група може перечитати багато даних. Ми налаштовуємо session.timeout.ms та max.poll.interval.ms так, щоб rebalance займав не більше 5 секунд, а через DLQ виключаємо застрягання на «поганих» повідомленнях.
Як ми гарантуємо доставку?
Згідно з документацією Apache Kafka, гарантії доставки визначаються комбінацією acks, ідемпотентності та транзакцій. Розглянемо ключові паттерни.
Ідемпотентний продюсер з acks=all
Фінансові транзакції, події замовлень — для будь-якого критичного потоку ми ставимо enable.idempotence=true та acks=all. Це виключає дублювання навіть при ретраях. Навантаження на пропускну здатність зростає незначно (5-10%), але надійність — на порядок.
// Java — ідемпотентний продюсер
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-1:9092,kafka-2:9092,kafka-3:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class);
props.put("schema.registry.url", "http://schema-registry:8081");
// Надійність
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
props.put(ProducerConfig.RETRIES_CONFIG, Integer.MAX_VALUE);
props.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 120_000);
// Продуктивність
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 65536); // 64KB батч
props.put(ProducerConfig.LINGER_MS_CONFIG, 5); // чекаємо до 5ms для наповнення батчу
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 67_108_864); // 64MB буфер
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
KafkaProducer<String, OrderEvent> producer = new KafkaProducer<>(props);
Відправка з обробкою помилок:
public CompletableFuture<RecordMetadata> sendOrderEvent(OrderEvent event) {
ProducerRecord<String, OrderEvent> record = new ProducerRecord<>(
"order-events",
event.getOrderId(),
event
);
CompletableFuture<RecordMetadata> future = new CompletableFuture<>();
producer.send(record, (metadata, exception) -> {
if (exception != null) {
if (exception instanceof RetriableException) {
log.error("Retriable error, Kafka will retry: {}", exception.getMessage());
} else {
log.error("Fatal producer error for order {}: {}", event.getOrderId(), exception.getMessage());
future.completeExceptionally(exception);
}
} else {
log.debug("Sent to partition {} offset {}", metadata.partition(), metadata.offset());
future.complete(metadata);
}
});
return future;
}
Приклад конфігурації консьюмера з ручним commit
Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-1:9092,kafka-2:9092,kafka-3:9092");
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "order-processor");
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
consumerProps.put("schema.registry.url", "http://schema-registry:8081");
consumerProps.put("specific.avro.reader", true);
// Вимикаємо auto-commit
consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
// Session timeout
consumerProps.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30_000);
consumerProps.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 10_000);
consumerProps.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 500);
consumerProps.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300_000);
KafkaConsumer<String, OrderEvent> consumer = new KafkaConsumer<>(consumerProps);
consumer.subscribe(List.of("order-events"), new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
commitCurrentOffsets();
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
log.info("Assigned partitions: {}", partitions);
}
});
Map<TopicPartition, OffsetAndMetadata> pendingOffsets = new HashMap<>();
try {
while (!shutdown.get()) {
ConsumerRecords<String, OrderEvent> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, OrderEvent> record : records) {
try {
processOrder(record.value());
pendingOffsets.put(
new TopicPartition(record.topic(), record.partition()),
new OffsetAndMetadata(record.offset() + 1)
);
} catch (NonRetriableException e) {
sendToDlq(record, e);
pendingOffsets.put(
new TopicPartition(record.topic(), record.partition()),
new OffsetAndMetadata(record.offset() + 1)
);
}
}
if (!pendingOffsets.isEmpty()) {
consumer.commitSync(pendingOffsets);
pendingOffsets.clear();
}
}
} finally {
consumer.close();
}
Чому ручне керування офсетами краще за auto-commit?
auto-commit створює ілюзію простоти, але приховує ризик: offset комітиться до обробки. При збої консьюмера між commit та обробкою повідомлення втрачається безповоротно. Ми завжди вимикаємо enable.auto.commit і комітимо вручну після успішної обробки запису. Це стандартна практика для продакшен-систем.
Порівняння режимів acks
| Режим | acks | Втрата даних при збої лідера | Продуктивність | Рекомендація |
|---|---|---|---|---|
| Fire-and-forget | 0 | Можлива | Максимальна | Тільки для неважливих логів |
| Лідер підтверджує | 1 | Можлива (неуспілі репліки) | Висока | Внутрішні сервіси з ретраями |
| Всі ISR | all (з min.insync.replicas) | Неможлива | Нижче на 10-15% | Критичні дані, замовлення, платежі |
Порівняння клієнтів Java та Python
| Критерій | Java (kafka-clients) | Python (confluent-kafka) |
|---|---|---|
| Продуктивність | Висока (нативна JVM) | Середня (через C-бенч) |
| Підтримка транзакцій | Повна | Обмежена |
| Керування офсетами | Ручне + async | В основному sync commit |
| Споживання пам'яті | Вище (JVM heap) | Нижче |
| Рекомендація | Критичні, high-throughput | Швидкі прототипи, середнє навантаження |
Як організувати паралельну обробку без втрати порядку?
Один потік poll + пул воркерів за ключем — класичний паттерн. Зберігаємо порядок подій одного замовлення, паралелимо між замовленнями.
Map<Integer, BlockingQueue<ConsumerRecord<String, OrderEvent>>> partitionQueues = new HashMap<>();
ExecutorService workers = Executors.newFixedThreadPool(12);
for (ConsumerRecord<String, OrderEvent> record : records) {
int partitionIndex = record.partition() % NUM_WORKERS;
workerQueues.get(partitionIndex).offer(record);
}
Процес розробки під ключ
- Аналітика — визначаємо топіки, ключі повідомлень, формат (Avro/JSON/Protobuf), групи консьюмерів.
- Реалізація — пишемо продюсерів з idempotence, консьюмерів з ручним commit, DLQ.
- Тестування — unit-тести з EmbeddedKafka, інтеграційні тести з Testcontainers, навантажувальне тестування.
- Деплой — налаштування моніторингу (consumer lag, помилки) та алертів.
Що входить в роботу
- Репозиторій з кодом (Java/Python) та документація по конфігурації.
- Docker-образи з налаштуваннями для staging/production.
- Інструкція по деплою та роллбеку.
- Навчання команди замовника (1-2 години).
- Гарантія на код — 3 місяці безкоштовної підтримки.
Строки
Від 6 до 10 робочих днів залежно від складності. Замовте оцінку вашого проєкту — відповімо протягом дня.
Типові помилки (чек-лист)
- [ ] Не включений enable.idempotence — дублікати при ретраях.
- [ ] Використовується auto-commit — втрата повідомлень при збої консьюмера.
- [ ] Занадто великий max.poll.records — довгий rebalance.
- [ ] Немає обробки NonRetriableException — втрата повідомлення без DLQ.
- [ ] Різні версії серіалізатора у продюсера та консьюмера — помилки десеріалізації.
Якщо ви впізнали в цьому свій проєкт — зв'яжіться з нами. Ми розробили продюсерів та консьюмерів для 15+ проєктів з навантаженням до 100 000 повідомлень на секунду.
Гарантія: фіксуємо строки в договорі, підтримуємо 3 місяці після здачі.







