【发布时间】:2020-05-10 04:55:34
【问题描述】:
我有一个要求,如果我可以删除内存中的 kafka 消息,因为我有 max-poll-records: 10. 所以场景是:在处理记录时,如果我的程序遇到任何错误,我不需要再处理存储在内存中的任何剩余记录。 例如:我一次获取 10 条记录,因为我的 max-poll-interval 为 10。我成功处理了 5 条记录(手动提交)但在第 6 条记录期间我遇到错误,现在我必须从内存中删除所有剩余的 5 条记录.下面是我的监听代码:
@KafkaListener(topics = "#{'${kafka.consumer.allTopicList}'.split(',')}", groupId = Constant.GROUP_ID)
public void consumeAllTopics(@Header(KafkaHeaders.RECEIVED_TOPIC) String topic, String message, Acknowledgment acknowledgment) {
switch (topic) {
case Constant.toipc1:
if (!StringUtils.isEmpty(message)) {
try {
//processing logic
acknowledgment.acknowledge();
}catch(Exception e){e.printStackTrace();}
}
我想通过代码删除记录。请帮助我了解是否可能,如果可以,我该如何实现。
【问题讨论】:
标签: spring-boot apache-kafka kafka-consumer-api spring-kafka