【问题标题】:org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method could not be invoked with the incoming messageorg.springframework.kafka.listener.ListenerExecutionFailedException:无法使用传入消息调用侦听器方法
【发布时间】:2019-02-25 11:30:25
【问题描述】:

我是 Apache Kafka 的新手,能够从发件人发送消息(JSON 格式),但无法在消费者服务中消费。

这是我的代码:

发件人服务

@Service
public class SenderService {
    private static final Logger LOG = LoggerFactory.getLogger(SenderService.class);
    
    @Autowired
    private KafkaTemplate<String, IdName> kafkaTemplate;
    
    @Value("${app.topic.email}")
    private String topic;
    
    public void send(IdName idName) {
        LOG.info("Sending Data='{}' to topic='{}' ", idName, topic);
        
        Message<IdName> message = MessageBuilder.withPayload(idName).setHeader(KafkaHeaders.TOPIC, topic)
            .setHeader(KafkaHeaders.MESSAGE_KEY, "TestMessage")
            .build();
        kafkaTemplate.send(message);
    }
}

###Consumer Service

@Service
public class ConsumerService {
    private static final Logger LOG = LoggerFactory.getLogger(ConsumerService.class);
    
    @Autowired
    MailSender mailSender;
    
    @KafkaListener(topics = "${app.topic.email}")
    public void receive(@Payload IdName data,
        @Header MessageHeaders headers) throws Exception{
        LOG.info("Received data='{}'", data);
    }
 }

我收到以下异常

2018-09-21 11:01:41.738 ERROR 63487 --- [ntainer#0-0-C-1] o.s.kafka.listener.LoggingErrorHandler   : Error while processing: ConsumerRecord(topic = emailclient, partition = 0, offset = 4, CreateTime = 1537507901660, serialized key size = 11, serialized value size = 122, headers = RecordHeaders(headers = [RecordHeader(key = __TypeId__, value = [99, 111, 109, 46, 97, 98, 99, 112, 108, 117, 115, 100, 46, 97, 112, 112, 108, 105, 99, 97, 116, 105, 111, 110, 46, 100, 111, 109, 97, 105, 110, 46, 73, 100, 78, 97, 109, 101])], isReadOnly = false), key = TestMessage, value = com.*****.****.domain.IdName@4a2d3006)

org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method could not be invoked with the incoming message
Endpoint handler details:
Method [public void com.*****.****.service.ConsumerService.receive(com.*****.****.domain.IdName,org.springframework.messaging.MessageHeaders) throws java.lang.Exception]

Bean [com.*****.****.service.ConsumerService@4649d70a]; nested exception is org.springframework.messaging.MessageHandlingException: Missing header 'headers' for method parameter type [class org.springframework.messaging.MessageHeaders], failedMessage=GenericMessage [payload=com.*****.****.domain.IdName@4a2d3006, headers={kafka_offset=4, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@53fdf4bb, kafka_timestampType=CREATE_TIME, kafka_receivedMessageKey=TestMessage, kafka_receivedPartitionId=0, kafka_receivedTopic=emailclient, kafka_receivedTimestamp=1537507901660, __TypeId__=[B@707343c2}]

谁能帮帮我?

【问题讨论】:

    标签: java apache-kafka


    【解决方案1】:

    如果您想访问侦听器中的所有标头,则使用了错误的注释,它应该是@Headers 而不是@Header。生成的代码如下所示:

    @Service
    public class ConsumerService {
    private static final Logger LOG = LoggerFactory.getLogger(ConsumerService.class);
    
        @Autowired
        MailSender mailSender;
    
        @KafkaListener(topics = "${app.topic.email}")
        public void receive(@Payload IdName data,
            @Headers MessageHeaders headers) throws Exception{
            LOG.info("Received data='{}'", data);
        }
     }
    

    您当然也可以只注入消息密钥,如下所示:

    @Service
    public class ConsumerService {
    private static final Logger LOG = LoggerFactory.getLogger(ConsumerService.class);
    
        @Autowired
        MailSender mailSender;
    
        @KafkaListener(topics = "${app.topic.email}")
        public void receive(@Payload IdName data,
            @Header(KafkaHeaders.MESSAGE_KEY) String messageKey) throws Exception{
            LOG.info("Received data='{}'", data);
        }
     }
    

    【讨论】:

      猜你喜欢
      • 2020-01-17
      • 2018-01-15
      • 2022-01-19
      • 1970-01-01
      • 2018-06-09
      • 2022-06-13
      • 2014-07-09
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多