【发布时间】:2021-05-21 14:02:57
【问题描述】:
我有一个使用 Spring Cloud Stream 和 Spring Kafka 的应用程序,它处理 Avro 消息。该应用程序运行良好,但现在我想添加一些错误处理。
目标:我想捕获反序列化异常,用异常详细信息 + 原始 Kafka 消息 + 自定义上下文信息构建一个新对象,并将此对象推送到专用 Kafka 主题。基本上是一个DLQ,但是原始消息会被截取和修饰。
问题:虽然我可以拦截异常,但我不知道如何从 Kafka 获取原始消息(TODO 1,下面)。我已经浏览过 ConsumerAwareErrorHandler.handle 中返回的数据对象,但我没有看到它。
下面是我的代码:
@EnableBinding(EventStream.class)
@SpringBootApplication
@Slf4j
public class SpringcloudApplication {
public static void main(String[] args) {
SpringApplication.run(SpringcloudApplication.class, args);
}
/* Configure custom exception handler */
@Bean
public ListenerContainerCustomizer<AbstractMessageListenerContainer<?, ?>> cust() {
return (container, destination, group) -> {
container.setErrorHandler(new ConsumerAwareErrorHandler() {
@Override
public void handle(Exception thrownException, ConsumerRecord<?, ?> data, Consumer<?, ?> consumer) {
log.info("Got error with data: {}", data);
// TODO 1 - How to get original message?
// TODO 2 - Send to dedicated (DLQ) topic
}
});
};
}
@StreamListener(EventStream.INBOUND)
public void consumeEvent(@Payload Message message) {
log.info("Consuming event --> {}", message.toString());
produceEvent(message);
}
@Autowired private EventStream eventStream;
public Boolean produceEvent(Message message) {
log.info("Producing event --> {}", message.toString());
return eventStream
.producer()
.send(MessageBuilder.withPayload(message)
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_JSON)
.build());
}
}
还有属性文件:
spring:
cloud:
stream:
default-binder: kafka
default:
consumer:
useNativeEncoding: true
producer:
useNativeEncoding: true
kafka:
binder:
brokers: localhost:9092
producer-properties:
key.serializer: org.apache.kafka.common.serialization.StringSerializer
value.serializer: io.confluent.kafka.serializers.KafkaAvroSerializer
schema.registry.url: "http://localhost:8081"
consumer-properties:
key.deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
value.deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
schema.registry.url: "http://localhost:8081"
specific.avro.reader: true
spring.deserializer.key.delegate.class: org.apache.kafka.common.serialization.StringDeserializer
spring.deserializer.value.delegate.class: io.confluent.kafka.serializers.KafkaAvroDeserializer
bindings:
event-consumer:
destination: data_stream_in # incoming topic
contentType: application/**avro
group: data_stream_consumer
event-producer:
destination: data_stream_out
contentType: application/**avro
我正在使用以下版本:
- Spring Boot 2.3.2.RELEASE
- Spring Cloud:Hoxton.SR8
- spring-cloud-stream-binder-kafka 3.0.8.RELEASE
- spring-kafka 2.5.12
感谢任何帮助!
【问题讨论】:
标签: avro spring-kafka spring-cloud-stream spring-cloud-stream-binder-kafka