📨 Раздел 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. Если бы брокер проверял фильтры, ему пришлось бы декомпрессировать сетевой пакет, распарсить каждое сообщение в Java-кучу и применить предикат. Это убило бы пропускную способность брокера на 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) передача трафика между 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?

Ответ:

  1. Zero-Copy архитектура: Брокер Kafka передает сжатые батчи напрямую из дискового Page Cache в сетевой сокет через системный вызов sendfile(2). Серверная фильтрация потребовала бы распаковки сжатых батчей в память JVM брокера и разбора формата сообщений, что превратило бы I/O-bound брокер в CPU-bound узкое горлышко.
  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 позволяет отсекать ненужные записи с нулевыми затратами на десериализацию.

Связанные темы