Як реалізувати фільтрацію повідомлень на стороні консьюмера
В 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 брокер не фільтрує дані?».
Причини мають фундаментальний характер:
- Збереження принципу Zero-Copy (
sendfile): Брокер передає стисненийRecordBatchіз дискового Page Cache безпосередньо в мережевий сокет мережевої карти без копіювання даних в User Space JVM. Якби брокер перевіряв фільтри, йому довелося б розпакувати мережевий пакет, розпарсити кожне повідомлення в купу (heap) JVM і застосувати предикат. Це знизило б пропускну здатність брокера на 95%. - Пасивний спільний лог для різнорідних споживачів:
В один і той самий топік
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?
Відповідь:
- Zero-Copy архітектура: Брокер Kafka передає стиснені пакети безпосередньо з дискового Page Cache у мережевий сокет за допомогою системного виклику
sendfile(2). Серверна фільтрація вимагала б розпакування стиснених батчів у пам’ять JVM брокера та парсингу формату повідомлень, що перетворило б I/O-bound брокер на вузьке горлечко, обмежене CPU. - Імутабельність партицій: У Kafka всі групи споживачів використовують один і той самий незмінний лог партиції. У чергах (JMS, RabbitMQ) кожне повідомлення адресується індивідуальним чергам підписників, тому брокер може легко застосувати фільтр на етапі маршрутизації (Exchange). У моделі спільного логу фільтрація на брокері зруйнувала б універсальність послідовних зміщень (
offset).
2. Що станеться, якщо в циклі консьюмера відфільтрувати повідомлення, але забути закомітити їхні зміщення?
Відповідь: Відбудеться каскадний повтор обробки відфільтрованих повідомлень (Replay Trap):
- Якщо консьюмер відфільтрував повідомлення з офсетами 100–199, а опрацював і закомітив лише офсет 200, то зміщення закріпиться успішно.
- Однак якщо надійшов цілий батч, що складається виключно з відфільтрованих повідомлень (офсети 300–400), і розробник помістив виклик
commitSync()лише всередину блоку успішної бізнес-обробки, зміщення 400 закомічено не буде! - При перезапуску застосунку або ребалансуванні консьюмер знову вичитає офсети 300–400, повторно витратить ресурси CPU на їхню фільтрацію і знову не зможе зсунути покажчик, ризикуючи застрягти в нескінченному циклі.
- Правило: Зміщення зобов’язане фіксуватися для всіх прочитаних записів батчу, незалежно від того, чи були вони відфільтровані, чи успішно оброблені.
3. Як працює параметр setAckDiscarded(true) у Spring Kafka ConcurrentKafkaListenerContainerFactory?
Відповідь:
Параметр setAckDiscarded керує підтвердженням зміщень тих записів, які були відхилені предикатом RecordFilterStrategy:
- Якщо
setAckDiscarded = true(за замовчуванням): відфільтрований запис не передається в метод@KafkaListener, але його зміщення негайно позначається як оброблене вAcknowledgmentконтейнера. Офсет просувається вперед у штатному режимі. - Якщо
setAckDiscarded = false: зміщення відфільтрованого повідомлення не підтверджується контейнером. Це має сенс лише при специфічних режимах ручного коміту (MANUAL_IMMEDIATE), інакше це може призвести до розсинхронізації позиції читання або зависання.
4. Як організувати фільтрацію за типом події з нульовими накладними витратами на CPU консьюмера?
Відповідь: Найбільш продуктивне рішення — фільтрація за байтовими заголовками (Record Headers) до десеріалізації:
- Продюсер додає бінарний заголовок:
record.headers().add("type", "PAYMENT".getBytes(StandardCharsets.UTF_8)). - Консьюмер налаштовується зі стандартним
ByteArrayDeserializer(або відкладеним лінивим десеріалізатором). - У циклі консьюмер перевіряє байти заголовка
record.headers().lastHeader("type"). - Якщо заголовок не відповідає очікуваному типу події, запис миттєво відкидається без створення Java-об’єктів у купі та без парсингу JSON/Avro/Protobuf.
- Навантаження на 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 дозволяє відсікати непотрібні записи з нульовими витратами на десеріалізацію.