You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

SpringBoot中如何在监听器级别基于Kafka Header值过滤消息?

问题:基于Kafka Header过滤消息时应用启动失败

我希望在监听器级别根据Kafka Header的值过滤消息,编写了自定义的KafkaConsumerConfig.java配置类,代码如下:

package com.example.test.consumer;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.header.Header;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;

import java.util.HashMap;
import java.util.Map;

@Configuration
public class KafkaConsumerConfig {

    @Value("${spring.kafka.consumer.group-id}")
    private String group_id;
    
    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrap_server;

    @Bean
    ConsumerFactory<String, String> consumerFactory() {
        Map<String, Object> config = new HashMap<>();
        config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap_server);
        config.put(ConsumerConfig.GROUP_ID_CONFIG, group_id);
        config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        return new DefaultKafkaConsumerFactory<>(config);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setRecordFilterStrategy(consumerRecord -> {
            if (consumerRecord.headers().lastHeader("CUSTOM_HEADER").key().equals("VALID")) {
                return false;
            }
            return true;
        });
        return factory;
    }
}

但应用启动失败,报错信息如下:

Cancelled in-flight API_VERSIONS request with correlation id 20 due to node -1 being disconnected


问题分析与解决

1. 启动连接错误排查

报错node -1 being disconnected本质是Kafka客户端无法连接指定的bootstrap服务器,常见原因:

  • 配置错误:spring.kafka.bootstrap-servers的地址/端口写错,或Kafka服务未启动
  • 网络限制:应用所在机器无法访问Kafka服务器(防火墙、网络策略拦截)
  • 注入异常:检查配置文件中spring.kafka.bootstrap-servers的键名是否正确,确保@Value能正常读取值

可先用Kafka命令行工具验证连接:

kafka-topics.sh --list --bootstrap-server <你的bootstrap地址>

2. 过滤逻辑的两处关键错误

即使连接问题解决,原代码的过滤逻辑存在致命问题:

  • 空指针风险:如果消息没有CUSTOM_HEADER,lastHeader("CUSTOM_HEADER")会返回null,调用key()会直接抛出空指针
  • 逻辑错误:Header.key()返回的是Header的名称(即CUSTOM_HEADER),而非Header存储的实际值,应该用Header.value()获取值后转成字符串比较

3. 修正后的完整配置类

package com.example.test.consumer;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.header.Header;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;

import java.util.HashMap;
import java.util.Map;

@Configuration
public class KafkaConsumerConfig {

    @Value("${spring.kafka.consumer.group-id}")
    private String groupId;
    
    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Bean
    ConsumerFactory<String, String> consumerFactory() {
        Map<String, Object> config = new HashMap<>();
        config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        config.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        return new DefaultKafkaConsumerFactory<>(config);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setRecordFilterStrategy(consumerRecord -> {
            Header customHeader = consumerRecord.headers().lastHeader("CUSTOM_HEADER");
            if (customHeader != null) {
                String headerValue = new String(customHeader.value());
                // 返回true=过滤消息,返回false=保留消息
                return !"VALID".equals(headerValue);
            }
            // 无目标Header时过滤消息
            return true;
        });
        return factory;
    }
}

内容的提问来源于stack exchange,提问作者Sayan Saha

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 06:20:35