Java中使用@KafkaListener消费Kafka消息时如何获取siebel-id参数?
提取方法
你可以通过Spring Kafka提供的@Header注解直接从消息中提取对应参数,无需修改现有消费者配置,分两种场景处理:
注意:需要导入的注解路径为
org.springframework.messaging.handler.annotation.Header、org.springframework.kafka.support.KafkaHeaders,避免导包错误。
场景1:siebel-id作为消息Key传递
你的消费者已经配置了String类型的Key反序列化器,直接注入消息Key即可:
@KafkaListener( containerFactory = "kafkaChangeClientPhoneListenerContainerFactory", topics = "${kafka.topic.changeClientPhone}" ) public void consume(ChangeClientPhoneEvent changeClientPhoneEvent, @Header(KafkaHeaders.RECEIVED_KEY) String siebelId) { // siebelId即为提取到的参数,可直接作为String类型使用 //TODO }
场景2:siebel-id作为自定义消息头传递
直接指定头参数的名称注入即可:
@KafkaListener( containerFactory = "kafkaChangeClientPhoneListenerContainerFactory", topics = "${kafka.topic.changeClientPhone}" ) public void consume(ChangeClientPhoneEvent changeClientPhoneEvent, @Header("siebel-id") String siebelId) { // 直接使用siebelId变量即可 //TODO }
兜底校验方案
如果你不确定siebel-id的存储位置,可以注入全量消息头遍历确认参数对应的key:
import org.springframework.messaging.handler.annotation.Headers; import java.util.Map; @KafkaListener( containerFactory = "kafkaChangeClientPhoneListenerContainerFactory", topics = "${kafka.topic.changeClientPhone}" ) public void consume(ChangeClientPhoneEvent changeClientPhoneEvent, @Headers Map<String, Object> headers) { // 遍历打印所有头参数,确认siebel-id的对应key headers.forEach((key, value) -> System.out.printf("头参数:%s = %s%n", key, value)); // 确认后提取参数 String siebelId = (String) headers.get("siebel-id"); // 若为消息Key则用以下代码提取:String siebelId = (String) headers.get(KafkaHeaders.RECEIVED_KEY); //TODO }
内容的提问来源于stack exchange,提问作者Maksym Rybalkin
相关产品推荐
相关产品推荐

