【发布时间】:2022-07-11 21:49:37
【问题描述】:
我使用 Spring Cloud Stream 和 Spring Cloud Function 创建了一个 Kafka Consumer,用于以批处理模式从 Kafka 主题消费消息。现在,我想将错误批次发送到死信队列,以便进一步调试错误。
我正在使用 Spring 重试在我的消费者方法中处理重试。但对于不可重试的异常,我希望将整个批次发送到 DLQ。
这就是我的消费者的样子:
@Bean
public Consumer<List<GenericRecord>> consume() {
return (message) -> {
processMessage(message);
}
}
这是错误处理配置的样子:
@Autowired
private DefaultErrorHandler errorHandler;
ListenerContainerCustomizer<AbstractMessageListenerContainer> c = new ListenerContainerCustomizer<AbstractMessageListenerContainer>() {
@Override
public void configure(AbstractMessageListenerContainer container, String destinationName, String group) {
container.setCommonErrorHandler(errorHandler);
}
}
使用 DeadRecordPublishinRecoverer 启用错误处理程序以将失败的消息发送到 DLQ:
@Bean
public DefaultErrorHandler errorHandler(KafkaOperations<String, Details> template) {
return new DefaultErrorHandler(new DeadLetterPublishingRecoverer(template,
(cr, e) -> new TopicPartition("error.topic.name", 0)),
new FixedBackOff(0, 0));
}
但这并没有向error.topic发送任何消息,并且从错误日志中我可以看到它正在尝试连接到localhost:9092,而不是我在spring.cloud.stream.kafka.binder.brokers中提到的代理。
如何配置 DLQ 提供程序以从 application.properties 读取 Kafka 元数据?
还有没有办法配置Supplier函数来创建DLQ提供者?
【问题讨论】:
标签: java spring spring-kafka spring-cloud-stream spring-cloud-function