You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.25 16:06:04