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

如何在失败消息进入死信主题前添加自定义Header?

解决方法

方案一:改用Consumer<Message<MyEventType>>手动处理失败与DLQ发送

将消费者参数从MyEventType改为Message<MyEventType>,获取完整消息头信息,在处理失败时手动构建带自定义Header的消息并发送到DLQ,同时关闭自动DLQ配置。

代码示例

  1. 修改消费者Bean:
@Bean
public Consumer<Message<MyEventType>> myConsumerFunction(StreamBridge streamBridge) {
    return message -> {
        try {
            MyEventType payload = message.getPayload();
            doStuff(payload);
        } catch (Exception e) {
            // 获取当前尝试次数,默认初始为1
            Integer attemptCount = message.getHeaders().get("x-processAttemptCount", Integer.class);
            attemptCount = attemptCount == null ? 1 : attemptCount + 1;

            // 构建新消息头,保留原有信息并添加自定义尝试次数
            Map<String, Object> newHeaders = new HashMap<>(message.getHeaders());
            newHeaders.put("x-processAttemptCount", attemptCount);
            MessageHeaders updatedHeaders = new MessageHeaders(newHeaders);

            // 生成新消息并发送到DLQ
            Message<MyEventType> updatedMessage = MessageBuilder.createMessage(message.getPayload(), updatedHeaders);
            streamBridge.send("my-dlt-name", updatedMessage);
        }
    };
}
  1. 修改配置关闭自动DLQ:
spring:
  cloud:
    stream:
      kafka:
        bindings:
          myConsumerFunction-in-0:
            consumer:
              enable-dlq: false # 关闭自动DLQ,改为手动发送

方案二:自定义DeadLetterPublishingRecoverer处理Header修改

结合Spring Retry与自定义死信发布恢复器,在消息重试失败后发送到DLQ前修改Header,贴合Spring Cloud Stream原生重试机制。

代码示例

  1. 配置重试模板与自定义恢复器:
@Bean
public RetryTemplate retryTemplate() {
    RetryTemplate retryTemplate = new RetryTemplate();
    // 配置最大重试次数(此处设置为2次,加上首次共3次尝试)
    SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy();
    retryPolicy.setMaxAttempts(2);
    retryTemplate.setRetryPolicy(retryPolicy);
    return retryTemplate;
}

@Bean
public DeadLetterPublishingRecoverer deadLetterPublishingRecoverer(KafkaTemplate<?, ?> kafkaTemplate) {
    return new DeadLetterPublishingRecoverer(kafkaTemplate, (record, exception) -> {
        // 获取并更新尝试次数
        Integer attemptCount = record.headers().lastHeader("x-processAttemptCount") != null
                ? Integer.parseInt(new String(record.headers().lastHeader("x-processAttemptCount").value()))
                : 1;
        attemptCount += 1;
        record.headers().add("x-processAttemptCount", String.valueOf(attemptCount).getBytes());

        // 指定死信主题
        return new TopicPartition("my-dlt-name", record.partition());
    });
}

@Bean
public Consumer<MyEventType> myConsumerFunction() {
    return value -> doStuff(value);
}
  1. 修改配置关联重试与恢复器:
spring:
  cloud:
    stream:
      kafka:
        bindings:
          myConsumerFunction-in-0:
            consumer:
              enable-dlq: false # 关闭自动DLQ,使用自定义恢复器
              retry:
                enabled: true
                template: retryTemplate # 关联自定义重试模板
              recoverer: deadLetterPublishingRecoverer # 关联自定义死信恢复器

死信主题消费逻辑

从DLQ消费时,读取x-processAttemptCountHeader判断是否超过阈值,决定发回原主题还是丢弃:

@Bean
public Consumer<Message<MyEventType>> dlqConsumerFunction(StreamBridge streamBridge) {
    return message -> {
        Integer attemptCount = message.getHeaders().get("x-processAttemptCount", Integer.class);
        int threshold = 3; // 设定最大尝试阈值

        if (attemptCount < threshold) {
            // 发回原主题重新处理
            streamBridge.send("myConsumerFunction-in-0", message.getPayload());
        } else {
            // 记录日志并丢弃消息
            log.error("消息处理超过阈值,已丢弃: {}", message.getPayload());
        }
    };
}

注意事项

  • MessageHeaders是不可变对象,必须通过MessageBuilder或构建新的MessageHeaders实例来修改Header,不能直接调用put()方法。
  • 使用StreamBridge发送消息时,需确保目标主题的绑定配置正确(自定义主题需在spring.cloud.stream.bindings中声明输出绑定)。

内容的提问来源于stack exchange,提问作者user3236794

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 18:52:54