Spring Kafka处理poison pill的替代方案:不发DLT直接存入数据库
实现方案
核心思路是替换Spring Kafka错误处理器中默认的死信队列发布恢复器,自定义消费失败恢复逻辑实现数据库写入:
步骤1:自定义消费失败恢复器
实现ConsumerRecordRecoverer接口,在接口方法中完成失败消息和异常信息的数据库写入逻辑:
import org.apache.commons.lang3.exception.ExceptionUtils; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.listener.ConsumerRecordRecoverer; import org.springframework.stereotype.Component; import java.time.LocalDateTime; import java.util.Objects; @Component public class DatabaseFailedMessageRecoverer implements ConsumerRecordRecoverer { // 注入你自己项目的数据库持久层实现(JPA Repository/MyBatis Mapper等) private final FailedMessageDao failedMessageDao; public DatabaseFailedMessageRecoverer(FailedMessageDao failedMessageDao) { this.failedMessageDao = failedMessageDao; } @Override public void accept(ConsumerRecord<?, ?> record, Exception exception) { // 构建失败消息存储实体,可按需扩展字段 FailedMessageEntity messageEntity = FailedMessageEntity.builder() .topic(record.topic()) .partition(record.partition()) .offset(record.offset()) .msgKey(Objects.toString(record.key(), null)) .msgContent(Objects.toString(record.value(), null)) .errorInfo(exception.getMessage()) .errorStack(ExceptionUtils.getStackTrace(exception)) .createTime(LocalDateTime.now()) .build(); // 写入数据库 failedMessageDao.insert(messageEntity); } }
步骤2:替换错误处理器配置
修改原有配置,将自定义恢复器注入到错误处理器中,替换原来的DeadLetterPublishingRecoverer:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.listener.SeekToCurrentErrorHandler; import org.springframework.util.backoff.FixedBackOff; @Configuration public class KafkaErrorConfig { @Bean public SeekToCurrentErrorHandler errorHandler(DatabaseFailedMessageRecoverer messageRecoverer) { // 第二个参数配置重试策略:此处为间隔1秒重试,最多重试3次,全部失败后走自定义恢复逻辑存数据库 return new SeekToCurrentErrorHandler(messageRecoverer, new FixedBackOff(1000L, 3)); } // 原有DeadLetterPublishingRecoverer相关的Bean可以直接删除 }
注意事项
- 如果你的Spring Kafka版本为2.8及以上,
SeekToCurrentErrorHandler已被标记为废弃,直接替换为DefaultErrorHandler即可,配置逻辑完全一致。 - 建议在自定义恢复器中添加数据库写入的异常捕获逻辑,避免数据库写入失败导致消费端无限重试同一条消息。
- 消息payload如果是自定义序列化对象,可直接在恢复器中强转后使用,无需统一转为字符串存储。
内容的提问来源于stack exchange,提问作者Balkrishna Kardam
相关产品推荐
相关产品推荐

