Как реализовать фильтрацию сообщений на стороне консьюмера
В 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. Если бы брокер проверял фильтры, ему пришлось бы декомпрессировать сетевой пакет, распарсить каждое сообщение в Java-кучу и применить предикат. Это убило бы пропускную способность брокера на 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) передача трафика между Availability Zones стоит реальных денег. Читать гигабайты ненужных данных ради 1% полезной нагрузки — прямые финансовые потери компании.
- Consumer Lag Skew: Потребитель тратит время сетевого I/O на прокачку миллионов «пустых» записей, провоцируя ложные срабатывания мониторинга по лагу.
Альтернатива: Фильтрация через 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-bound узкое горлышко. - Иммутабельность партиций: В 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 позволяет отсекать ненужные записи с нулевыми затратами на десериализацию.