【发布时间】:2022-01-09 14:02:29
【问题描述】:
我正在尝试使用启用批处理模式的 Spring Cloud Stream 实现 DLQ
@Bean
public ListenerContainerCustomizer<AbstractMessageListenerContainer<?, ?>> customizer(BatchErrorHandler handler) {
return ((container, destinationName, group) -> {
if(dlqEnabledTopic.contains(destinationName))
container.setBatchErrorHandler(handler);});
}
@Bean
public BatchErrorHandler batchErrorHandler(KafkaOperations<String, byte[]> kafkaOperations) {
CustomDeadLetterPublishingRecoverer recoverer = new CustomDeadLetterPublishingRecoverer(kafkaOperations,
(cr, e) -> new TopicPartition(cr.topic()+"_dlq", cr.partition()));
return new RecoveringBatchErrorHandler(recoverer, new FixedBackOff(1000, 1));
}
但有一些疑问:
-
如何使用属性配置键/值序列化器 - 我的消息是字符串类型,但 KafkaOperations 使用的是 ByteArraySerializer
-
在批处理中有多条消息,但如果第一条消息失败,它会转到 DLQ 但看不到下一条消息的处理。
要求 - 在任何索引处,如果批处理失败,我只需将该消息发送到 DLQ,其余消息应再次处理。
- 现在批处理模式支持 DLQ 吗?就像记录模式一样,它可以使用属性启用
【问题讨论】:
标签: spring-boot spring-cloud spring-kafka spring-cloud-stream