📨 Розділ 15 · Питання #30

Як реалізувати фільтрацію повідомлень на стороні консьюмера

В Apache Kafka брокер не виконує фільтрацію вмісту повідомлень. Брокер спроєктований як високопродуктивне пасивне сховище послідовних байтів («розумні клієнти, простий брокер»)....


🟢 Junior Level

Чому фільтрація виконується на клієнті

В Apache Kafka брокер не виконує фільтрацію вмісту повідомлень. Брокер спроєктований як високопродуктивне пасивне сховище послідовних байтів («розумні клієнти, простий брокер»). Брокер не знає і не повинен знати бізнес-схему даних (JSON, Avro, Protobuf) та бізнес-правила мікросервісів.

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

[Продюсер] ──(Замовлення: $10, $500, $20, $1500)──> [Kafka Брокер (Топік: orders)]
                                                             │
                                                             ▼ (Мережа: передаються ВСІ записи)
                                                 [Consumer (Java)]
                                                             │
                                         ┌───────────────────┴───────────────────┐
                                         ▼                                       ▼
                             Amount > $100: ОБРОБИТИ                  Amount <= $100: ВІДХИЛИТИ
                             (Зберегти в БД, викликати API)           (Зміщення все одно комітиться!)

Базовий приклад на чистому Java Kafka Consumer

Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "vip-order-service");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(List.of("orders"));

while (running) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    
    for (ConsumerRecord<String, String> record : records) {
        // Фільтрація на стороні консьюмера
        if (isVipOrder(record.value())) {
            processVipOrder(record);
        }
        // Відфільтровані повідомлення пропускаються, але їхній офсет обов'язково враховується!
    }
    
    // Комітимо зміщення всього батчу: відфільтровані повідомлення НЕ будуть перечитуватися знову
    consumer.commitSync();
}

🟡 Middle Level

Фільтрація в Spring Kafka: RecordFilterStrategy

У застосунках на Spring Boot ручне написання розгалужень if-else всередині бізнес-методів вважається антипатерном. Spring Kafka надає декларативний інтерфейс RecordFilterStrategy:

@Configuration
public class KafkaConsumerConfig {

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, OrderDto> vipOrderListenerFactory(
            ConsumerFactory<String, OrderDto> consumerFactory) {

        ConcurrentKafkaListenerContainerFactory<String, OrderDto> factory = 
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);

        // Налаштування фільтра: повертає true, якщо повідомлення потрібно ВІДХИЛИТИ (discard)
        factory.setRecordFilterStrategy(record -> {
            OrderDto order = record.value();
            return order == null || order.amount() < 1000.0; // Відкидаємо не-VIP замовлення
        });

        // Налаштування підтвердження: відфільтровані повідомлення комітяться автоматично
        factory.setAckDiscarded(true);

        return factory;
    }
}

Використання в бізнес-сервісі:

@Service
public class VipOrderProcessor {

    // Сюди потраплять ТІЛЬКИ замовлення із сумою >= 1000.0
    @KafkaListener(topics = "orders", containerFactory = "vipOrderListenerFactory")
    public void handleVipOrder(OrderDto order) {
        log.info("Обробка VIP-замовлення: {}", order.id());
    }
}

Оптимізація: Фільтрація за заголовками (Header-based Filtering)

Десеріалізація важкого тіла повідомлення (JSON / XML) у Java-об’єкти споживає значну кількість тактів CPU та навантажує Garbage Collector алокаціями в пам’яті. Якщо в топіку відфільтровується 80–90% записів, десеріалізувати record.value() для кожного повідомлення вкрай неефективно.

Рішення: Продюсер виносить метадані маршрутизації в Kafka Headers:

// Перевірка заголовка БЕЗ десеріалізації тіла повідомлення
Header typeHeader = record.headers().lastHeader("eventType");
if (typeHeader != null && "VIP_PAYMENT".equals(new String(typeHeader.value(), StandardCharsets.UTF_8))) {
    // Лише тепер десеріалізуємо важкий JSON
    OrderDto order = objectMapper.readValue(record.value(), OrderDto.class);
    process(order);
}

🔴 Senior Level

Чому Kafka архітектурно відмовилася від Server-side фільтрації

Розробники часто запитують: «Чому в RabbitMQ є Routing Keys, в ActiveMQ є JMS SQL92 селектори, а в Kafka брокер не фільтрує дані?».

Причини мають фундаментальний характер:

  1. Збереження принципу Zero-Copy (sendfile): Брокер передає стиснений RecordBatch із дискового Page Cache безпосередньо в мережевий сокет мережевої карти без копіювання даних в User Space JVM. Якби брокер перевіряв фільтри, йому довелося б розпакувати мережевий пакет, розпарсити кожне повідомлення в купу (heap) JVM і застосувати предикат. Це знизило б пропускну здатність брокера на 95%.
  2. Пасивний спільний лог для різнорідних споживачів: В один і той самий топік orders дивляться різні консьюмер-групи:
    • Сервіс доставки (потрібні замовлення з фізичними товарами).
    • Сервіс цифрових ліцензій (потрібні виключно цифрові товари).
    • Бухгалтерія (потрібні абсолютно всі замовлення для фінансового звіту). Якби брокер фільтрував лог, йому довелося б підтримувати окремі індекси та копії даних під вимоги кожної окремої групи.

Архітектурний антипатерн «High-Ratio Consumer Filtering»

Якщо консьюмер вичитує з топіка 100 000 RPS, а після фільтрації залишає лише 1 000 RPS (99% відсіву), це є грубою помилкою проєктування архітектури:

НЕПРАВИЛЬНО (Consumer Filtering 99%):
[Продюсер] ──100K msg/sec──> [Topic: all-events] ──Мережа: 100 MB/sec──> [Consumer (Фільтрує 99%)]
                                                                        (Марнотратство мережі та CPU)

ПРАВИЛЬНО (Producer Content-Based Routing):
                             ┌───> [Topic: telemetry-events] ─────────> [Telemetry Consumer]
[Продюсер] ──Роутинг за типом┼───> [Topic: billing-events]   ─────────> [Billing Consumer]
                             └───> [Topic: audit-events]     ─────────> [Audit Consumer]

Приховані витрати надмірної фільтрації на клієнті:

  • Cloud Egress Costs: У хмарних провайдерах (AWS, GCP, Azure) міжзональний трафік (між Availability Zones) тарифікується окремо. Читати гігабайти зайвих даних заради 1% корисного навантаження — це прямі фінансові збитки.
  • Consumer Lag Skew: Споживач витрачає час введення-виведення на прокачування мільйонів «порожніх» записів, що провокує хибні спрацьовування алертингу щодо зростання лагу.

Альтернатива: Фільтрація через Kafka Streams / ksqlDB

Якщо продюсер належить зовнішній системі й не може маршрутизувати дані в різні топіки на етапі відправки, між вихідним топіком і цільовим сервісом впроваджується топологія потокової фільтрації:

StreamsBuilder builder = new StreamsBuilder();
KStream<String, OrderDto> sourceStream = builder.stream("orders");

// Потоковий предикат фільтрує дані та записує результат у легковаговий топік
sourceStream
    .filter((key, order) -> order.isVip())
    .to("orders-vip", Produced.with(Serdes.String(), new JsonSerde<>(OrderDto.class)));

// Мікросервіси підписуються вже на чистий відфільтрований топік "orders-vip"

4 Tricky Questions

1. Чому в Kafka брокер не підтримує server-side фільтрацію, на відміну від селекторів повідомлень у JMS або RabbitMQ?

Відповідь:

  1. Zero-Copy архітектура: Брокер Kafka передає стиснені пакети безпосередньо з дискового Page Cache у мережевий сокет за допомогою системного виклику sendfile(2). Серверна фільтрація вимагала б розпакування стиснених батчів у пам’ять JVM брокера та парсингу формату повідомлень, що перетворило б I/O-bound брокер на вузьке горлечко, обмежене CPU.
  2. Імутабельність партицій: У Kafka всі групи споживачів використовують один і той самий незмінний лог партиції. У чергах (JMS, RabbitMQ) кожне повідомлення адресується індивідуальним чергам підписників, тому брокер може легко застосувати фільтр на етапі маршрутизації (Exchange). У моделі спільного логу фільтрація на брокері зруйнувала б універсальність послідовних зміщень (offset).

2. Що станеться, якщо в циклі консьюмера відфільтрувати повідомлення, але забути закомітити їхні зміщення?

Відповідь: Відбудеться каскадний повтор обробки відфільтрованих повідомлень (Replay Trap):

  1. Якщо консьюмер відфільтрував повідомлення з офсетами 100–199, а опрацював і закомітив лише офсет 200, то зміщення закріпиться успішно.
  2. Однак якщо надійшов цілий батч, що складається виключно з відфільтрованих повідомлень (офсети 300–400), і розробник помістив виклик commitSync() лише всередину блоку успішної бізнес-обробки, зміщення 400 закомічено не буде!
  3. При перезапуску застосунку або ребалансуванні консьюмер знову вичитає офсети 300–400, повторно витратить ресурси CPU на їхню фільтрацію і знову не зможе зсунути покажчик, ризикуючи застрягти в нескінченному циклі.
  4. Правило: Зміщення зобов’язане фіксуватися для всіх прочитаних записів батчу, незалежно від того, чи були вони відфільтровані, чи успішно оброблені.

3. Як працює параметр setAckDiscarded(true) у Spring Kafka ConcurrentKafkaListenerContainerFactory?

Відповідь: Параметр setAckDiscarded керує підтвердженням зміщень тих записів, які були відхилені предикатом RecordFilterStrategy:

  • Якщо setAckDiscarded = true (за замовчуванням): відфільтрований запис не передається в метод @KafkaListener, але його зміщення негайно позначається як оброблене в Acknowledgment контейнера. Офсет просувається вперед у штатному режимі.
  • Якщо setAckDiscarded = false: зміщення відфільтрованого повідомлення не підтверджується контейнером. Це має сенс лише при специфічних режимах ручного коміту (MANUAL_IMMEDIATE), інакше це може призвести до розсинхронізації позиції читання або зависання.

4. Як організувати фільтрацію за типом події з нульовими накладними витратами на CPU консьюмера?

Відповідь: Найбільш продуктивне рішення — фільтрація за байтовими заголовками (Record Headers) до десеріалізації:

  1. Продюсер додає бінарний заголовок: record.headers().add("type", "PAYMENT".getBytes(StandardCharsets.UTF_8)).
  2. Консьюмер налаштовується зі стандартним ByteArrayDeserializer (або відкладеним лінивим десеріалізатором).
  3. У циклі консьюмер перевіряє байти заголовка record.headers().lastHeader("type").
  4. Якщо заголовок не відповідає очікуваному типу події, запис миттєво відкидається без створення Java-об’єктів у купі та без парсингу JSON/Avro/Protobuf.
  5. Навантаження на CPU та паузи складання сміття (GC) знижуються в 10–20 разів порівняно з повною десеріалізацією кожного повідомлення.

🎯 Шпаргалка для інтерв’ю

30-секундна відповідь (Elevator Pitch)

«У Kafka фільтрація повідомлень принципово виконується на стороні Consumer, оскільки брокер спроєктований як пасивне сховище і використовує Zero-Copy передачу байтів без парсингу вмісту. У Spring Boot фільтрація елегантно реалізується через RecordFilterStrategy у фабриці слухачів з автоматичним підтвердженням відкинутих зміщень (ackDiscarded=true). Для оптимізації CPU фільтрацію слід проводити за заголовками (Record Headers) до етапу десеріалізації тіла повідомлення. Якщо консьюмер відсіює переважну більшість даних (>80%), це свідчить про архітектурну помилку: у таких випадках правильніше використовувати маршрутизацію на стороні продюсера за окремими топіками або проміжний фільтр на Kafka Streams».

Ключові метрики для Production-моніторингу

  • kafka_consumer_records_consumed_total: загальна кількість вичитаних повідомлень із мережі.
  • kafka_consumer_records_filtered_total: кількість повідомлень, відкинутих фільтром консьюмера.
  • filtered_ratio = records_filtered / records_consumed: частка відфільтрованих даних. Якщо вона стабільно перевищує 0.8 — потрібен поділ топіків.

Червоні прапорці (Чого говорити НЕ МОЖНА)

  • ❌ «У Kafka брокер сам відфільтрує непотрібні повідомлення, якщо правильно налаштувати топік» — Брокер Kafka не фільтрує корисне навантаження; вся фільтрація відбувається виключно на клієнті.
  • ❌ «Якщо повідомлення відфільтроване, його зміщення комітити не потрібно» — Незакомічений офсет призведе до повторного читання тих самих записів і зациклення при рестарті.
  • ❌ «Фільтрувати 99% повідомлень у мікросервісі — це стандартна нормальна практика» — Це архітектурний антипатерн, що призводить до невиправданих витрат на мережевий трафік і перевантаження CPU.
  • ❌ «Для фільтрації обов’язково десеріалізувати все тіло повідомлення в DTO» — Фільтрація за Kafka Headers дозволяє відсікати непотрібні записи з нульовими витратами на десеріалізацію.

Пов’язані теми