【发布时间】:2022-01-13 06:05:35
【问题描述】:
我正在处理错误处理实施,并有一个问题。
让我解释一下问题:
我收到一批消息,我在 for 循环中对每个消息进行数据库查找,然后我需要收集列表中所有查找的对象,并使用此对象列表调用批量插入存储过程和批量更新存储过程。
现在让我们假设在查找过程中发生了一些异常。我想重试这条消息。对于这种情况,我尝试使用DefaultErrorHandler。但是有一个问题,根据文档,当抛出带有元素索引的 BatchListenerFailedException 时,它会提交索引之前记录的偏移量。但正如我所说,我需要在查找后执行批量插入和更新,所以我不想在索引之前提交偏移量,这些记录还没有插入/更新到数据库中。
这是否意味着我唯一的选择是使用RetryingBatchErrorHandler 每次都重试整个批次?我可以以某种方式继续处理不产生错误的消息吗?
另外,如果RetryingBatchErrorHandler 是唯一的选择,我怎么能确定在较长的退避期(指数退避)的情况下,kafka 不会杀死我的消费者并且不会启动重新平衡?
我目前的实现:
RetryingBatchErrorHandler retryingBatchErrorHandler =
new RetryingBatchErrorHandler(backoff,
(consumerRecord, e) ->
log.error("Backoff attempts exhausted for the record with offset={}, partition={}, value={}, offset committed.",
consumerRecord.offset(), consumerRecord.partition(), consumerRecord.value()));
factory.setBatchErrorHandler(retryingBatchErrorHandler);
更新:请参阅 Artem 答案中的 cmets。
这就是如何将查找步骤包装到 retryTemplate
LookedUpRequest lookedUpRequest = retryTemplate.execute(ctx -> {
//Lookup step
return lookup.process(request);
});
如果它失败了,那么它将进一步为批处理错误处理程序抛出异常,RetryingBatchErrorHandler 根据其策略重试批处理
【问题讨论】:
标签: java spring spring-boot apache-kafka spring-kafka