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 месяца после сдачи.







