【发布时间】:2021-09-15 02:57:21
【问题描述】:
我在 kafka 中使用批处理,在带有批处理模式的 spring cloud stream kafka binder 中不支持重试,有一个选项可以配置 SeekToCurrentBatchErrorHandler(使用 ListenerContainerCustomizer)来实现类似的功能,以便在 binder 中重试.
我尝试了同样的方法,但使用了 SeekToCurrentBatchErrorHandler,但它重试的时间超过了设置的 3 次。
-
我该怎么做? 我想重试整个批次。
-
如何将整个批次发送到 dlq 主题?就像我曾经将deliveryAttempt(retry)匹配到3然后发送到DLQ主题的记录监听器一样,签入监听器。
我已经检查了this link, which is exactly my issue,但是一个例子会很有帮助,使用这个库 spring-cloud-stream-kafka-binder,我可以实现这一点。请举例说明,我是新手。
目前我有以下代码。
@Configuration
public class ConsumerConfig {
@Bean
public ListenerContainerCustomizer<AbstractMessageListenerContainer<?, ?>> customizer() {
return (container, dest, group) -> {
container.getContainerProperties().setAckOnError(false);
SeekToCurrentBatchErrorHandler seekToCurrentBatchErrorHandler
= new SeekToCurrentBatchErrorHandler();
seekToCurrentBatchErrorHandler.setBackOff(new FixedBackOff(0L, 2L));
container.setBatchErrorHandler(seekToCurrentBatchErrorHandler);
//container.setBatchErrorHandler(new BatchLoggingErrorHandler());
};
}
}
听者:
@StreamListener(ActivityChannel.INPUT_CHANNEL)
public void handleActivity(List<Message<Event>> messages,
@Header(name = KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment
acknowledgment,
@Header(name = "deliveryAttempt", defaultValue = "1") int
deliveryAttempt) {
try {
log.info("Received activity message with message length {}", messages.size());
nodeConfigActivityBatchProcessor.processNodeConfigActivity(messages);
acknowledgment.acknowledge();
log.debug("Processed activity message {} successfully!!", messages.size());
} catch (MessagePublishException e) {
if (deliveryAttempt == 3) {
log.error(
String.format("Exception occurred, sending the message=%s to DLQ due to: ",
"message"),
e);
publisher.publishToDlq(EventType.UPDATE_FAILED, "message", e.getMessage());
} else {
throw e;
}
}
}
看到@Gary 的回复后,添加了带有RetryingBatchErrorHandler 的ListenerContainerCustomizer @Bean,但无法导入该类。附上截图。
【问题讨论】:
标签: java spring apache-kafka spring-cloud-stream spring-cloud-stream-binder-kafka