Розробка продюсерів та консьюмерів Kafka під ключ

Наша компанія займається розробкою, підтримкою та обслуговуванням сайтів будь-якої складності. Від простих односторінкових сайтів до масштабних кластерних систем, побудованих на мікро сервісах. Досвід розробників підтверджено сертифікатами від вендорів.

Розробка та обслуговування будь-яких видів сайтів:

Інформаційні сайти або веб-програми
Сайти візитки, landing page, корпоративні сайти, онлайн каталоги, квіз, промо-сайти, блоги, ресурси новин, інформаційні портали, форуми, агрегатори
Сайти або веб-програми електронної комерції
Інтернет-магазини, B2B-портали, маркетплейси, онлайн-обмінники, кешбек-сайти, біржі, дропшиппінг-платформи, парсери товарів
Веб-програми для управління бізнес-процесами
CRM-системи, ERP-системи, корпоративні портали, системи управління виробництвом, парсери інформації
Сайти або веб-програми електронних послуг
Дошки оголошень, онлайн-школи, онлайн-кінотеатри, конструктори сайтів, портали надання електронних послуг, відеохостинги, тематичні портали

Це лише деякі з технічних типів сайтів, з якими ми працюємо, і кожен із них може мати свої специфічні особливості та функціональність, а також бути адаптованим під конкретні потреби та цілі клієнта.

Послуги, які ми пропонуємо
Показано 1 з 1Усі 2062 послуг
Розробка продюсерів та консьюмерів Kafka під ключ
Складний
~3-5 днів
Часті запитання

Наші компетенції:

Етапи розробки

Останні роботи

  • image_website-b2b-advance_0.webp
    Розробка сайту компанії B2B ADVANCE
    1359
  • image_web-applications_feedme_466_0.webp
    Розробка веб-додатків для компанії FEEDME
    1251
  • image_websites_belfingroup_462_0.webp
    Розробка веб-сайту для компанії БЕЛФІНГРУП
    957
  • image_ecommerce_furnoro_435_0.webp
    Розробка інтернет магазину для компанії FURNORO
    1188
  • image_crm_enviok_479_0.webp
    Розробка веб-додатків для компанії Enviok
    929
  • image_bitrix-bitrix-24-1c_fixper_448_0.webp
    Розробка веб-сайту для компанії ФІКСПЕР
    947

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);
}

Процес розробки під ключ

  1. Аналітика — визначаємо топіки, ключі повідомлень, формат (Avro/JSON/Protobuf), групи консьюмерів.
  2. Реалізація — пишемо продюсерів з idempotence, консьюмерів з ручним commit, DLQ.
  3. Тестування — unit-тести з EmbeddedKafka, інтеграційні тести з Testcontainers, навантажувальне тестування.
  4. Деплой — налаштування моніторингу (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 місяці після здачі.

Послуги бекенд-розробки: production-grade надійність

На production-сервері о 3:14 ночі черга Laravel Jobs перестала оброблятися — 40 000 необроблених завдань у Redis. Причина: worker упав через memory leak у статичній змінній Eloquent observer, supervisor не перезапустив через misconfigured stopwaitsecs. Ми розбирали такий інцидент на проекті з 500 RPS: діагностика 4 години, фікс — 20 хвилин. Щоб ви не втрачали гроші, пропонуємо послуги бекенд-розробки з акцентом на production-grade надійність — 10+ років досвіду, 50+ проектів, 5 років на ринку. Оцінимо ваш проект за 2 дні.

Які проблеми вирішуємо

N+1 запити: головний вбивця швидкості

N+1 — найпоширеніша причина повільних сторінок у Laravel-додатках. Стандартна історія: сторінка працювала нормально на dev з 10 записами, на production з 10 000 — 8-секундне завантаження.

Laravel Debugbar у dev-оточенні показує кількість запитів. Більше 20 — сигнал для audit.

Model::preventLazyLoading(! app()->isProduction());

Telescope для профілювання: логує всі запити, jobs, mail, notifications з деталізацією. Після впровадження eager loading час завантаження сторінки падає з 8 с до 0.3 с — у 27 разів.

Memory leak у статичних змінних

У Laravel Octane або Swoole додаток тримається в пам’яті між запитами. Статичні змінні не скидаються — призводять до неконтрольованого росту пам’яті. Використовуємо defer-функції та контейнерні біндинги для коректного скидання стану.

Неправильний connection pool

Rails, Laravel, Django відкривають нове з'єднання PostgreSQL на кожен PHP/Python процес. 100 воркерів — 100 з'єднань. PostgreSQL деградує від 200+ активних з'єднань через overhead на управління.

PgBouncer у transaction pooling: 1000 воркерів → 20–50 реальних з'єднань. Це знижує latency на 40% та зменшує витрати на хостинг на 30% — при середній вартості хостингу $2,000/міс економить $600/міс. GIN-індекс для JSONB до 100 разів швидший за B-tree при пошуку.

Як Octane справляється з високим навантаженням?

Laravel Octane (RoadRunner або Swoole) прибирає overhead bootstrap на кожен HTTP-запит. Приріст: 3–8x на синтетичних бенчмарках, 2–4x на реальних додатках. Важливо: не зберігати стан у статичних змінних — застосовуємо це на проектах >1000 RPS.

Як PostgreSQL допомагає уникнути повільних запитів?

Використовуємо composite indexes для WHERE + ORDER BY, partial indexes для фільтрів з високою селективністю, GIN-індекси для JSONB та full-text search. to_tsvector + GIN замість LIKE '%query%' — запобігає seq scan навіть на мільйонах записів. Аналізуємо плани через EXPLAIN ANALYZE та pg_stat_statements.

Як обрати стек для вашого проекту?

Стек Коли використовувати
Laravel + Octane CRUD, бізнес-логіка, REST/GraphQL API, адмінки
Node.js (Fastify) Realtime WebSocket, streaming, serverless, висока I/O concurrency
Go Високонавантажені мікросервіси (>10k RPS), gRPC, DevOps-інструменти
Django + DRF ML-пайплайни, інтеграція з AI, складна обробка даних
Ruby on Rails Швидкий MVP з багатим екосистемою гемів

Node.js виправданий для realtime: Laravel публікує події в Redis Pub/Sub, Node.js підписується та транслює клієнтам. Go — для goroutines (10k з'єднань на сервер — норма), але розробка повільніша, ніж Laravel.

Чому Redis критичний для продуктивності?

Redis виконує кілька ролей:

Роль Деталі
Кеш Кешування результатів важких запитів, фрагментів HTML
Черги Backend для Laravel Queue / Celery
Session store Distributed sessions в multi-instance оточенні
Pub/Sub Realtime події між сервісами
Rate limiting Sliding window counters для API throttling
Leaderboards Sorted Sets для рейтингів

Redis Cluster для горизонтального масштабування, Sentinel для автоматичного failover. Замовте консультацію щодо оптимізації Redis для вашого проекту.

Що входить в роботу під ключ

  • Архітектурне проектування (документація API, схема БД, діаграма сервісів)
  • Реалізація за узгодженим ТЗ з code review
  • Налаштування CI/CD (GitHub Actions, Docker), моніторингу (Sentry, Grafana), алертингу
  • Навантажувальне тестування (k6, wrk) зі звітом
  • Передача вихідних кодів, доступів, інструкція з деплою
  • Навчання команди замовника (2–3 сесії)
  • Гарантійна підтримка 1 місяць після здачі

Орієнтири по термінах

Задача Термін
REST API для мобільного/SPA (середня складність) 6–12 тижнів
Backend зі складною бізнес-логікою + інтеграції 12–20 тижнів
Високонавантажений сервіс на Go 8–16 тижнів
Міграція legacy PHP на Laravel 16–32 тижні

Вартість розраховується індивідуально після аналізу вимог до навантаження, інтеграцій та бізнес-логіки. Зв'яжіться з нами для безкоштовного аудиту вашого поточного backend — отримайте план оптимізації за 2 дні. Замовте консультацію та дізнайтеся, як знизити витрати на інфраструктуру на 30% без втрати продуктивності.