【发布时间】:2019-12-19 14:41:12
【问题描述】:
我创建了一个 Springboot 应用程序来将消息推送到 Kafka 主题。该应用程序运行良好。我正在尝试的是在将消息发送到 Kafka 主题时出现故障时处理异常。我在发送消息时使用错误通道来跟踪错误。但实际问题是,我能够看到错误消息,但我无法看到错误消息中失败的实际有效负载。实际上,我想记录该有效负载。
我尝试发送的 JSON 消息:{"key1":"value1"}
服务类:
@AllArgsConstructor
@EnableBinding(Source.class)
public class SendMessageToKafka {
private final Source source;
public void sendMessage(String sampleMessage) {
source.output.send(MessageBuilder.withPayLoad(sampleMessage).build());
}
@ServiceActivator(inputChannel = "errorChannel")
public void errorHandler(ErrorMessage em) {
System.out.println(em);
}
}
application.yml:
spring:
cloud:
stream:
bindings:
output:
producer:
error-channel-enabled: true
通过以上配置,当Kafka服务器宕机时,控制权来到errorHandler方法并打印消息。但我无法从错误消息中看到 {"key1":"value1"} 的实际有效负载。如何从错误消息中检索它?
【问题讨论】:
-
以上实现是基于来自[link] (github.com/spring-cloud/spring-cloud-stream/issues/795)的cmets
-
实现基于上述链接讨论中 Gary Russell 的评论
标签: java spring spring-boot spring-kafka spring-cloud-stream