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

Spring Cloud RabbitMQ如何拦截DLQ重发消息移除x-exception-stacktrace头

问题描述

我希望在消息重试次数耗尽、被重发至死信队列(DLQ)时拦截消息,最终目标是移除这类消息中的x-exception-stacktrace请求头。

当前配置

YAML配置

spring:
  application:
    name: sandbox
  cloud:
    function:
      definition: rabbitTest1Input
    stream:
      binders:
        rabbitTestBinder1:
          type: rabbit
          environment:
            spring:
              rabbitmq:
                addresses: localhost:55015
                username: guest
                password: guest
                virtual-host: test
 
      bindings:
        rabbitTest1Input-in-0:
          binder: rabbitTestBinder1
          consumer:
            max-attempts: 3
          destination: ex1
          group: q1
      rabbit:
        bindings:
          rabbitTest1Input-in-0:
            consumer:
              autoBindDlq: true
              bind-queue: true
              binding-routing-key: q1key
              deadLetterExchange: ex1-DLX
              dlqDeadLetterExchange: ex1
              dlqDeadLetterRoutingKey: q1key_dlq
              dlqTtl: 180000
              prefetch: 5
              queue-name-group-only: true
              republishToDlq: true
              requeueRejected: false
              ttl: 86400000

消费者代码(Kotlin)

@Configuration
class ConsumerConfig {

    companion object : KLogging()

    @Bean
    fun rabbitTest1Input(): Consumer<Message<String>> {
        return Consumer {
            logger.info("Received from test1 queue: ${it.payload}")
            throw AmqpRejectAndDontRequeueException("FAILED")  // 触发重试耗尽后重发至DLQ的逻辑
        }
    }
}

已尝试的方案及问题

  • 注册@GlobalChannelInterceptor方案
    该方案无法生效,原因是RabbitMessageChannelBinder重发消息到DLQ时使用的是自身私有的RabbitTemplate实例,不是Spring容器中自动注入的实例,拦截逻辑无法覆盖到这个私有模板的发送流程。

  • 继承RabbitMessageChannelBinder自定义Binder方案
    思路是继承原生Binder类,移除其中写入x-exception-stacktrace请求头的逻辑,再将自定义类注册为Bean替换原生实现,核心代码如下:

    /**
     * 基于RabbitMessageChannelBinder修改,目的是移除重发到DLQ的消息中携带的x-exception-stacktrace头
     */
    class RabbitMessageChannelBinderWithNoStacktraceRepublished 
        : RabbitMessageChannelBinder(...)
    
    @Configuration
    @Import(
        RabbitAutoConfiguration::class,
        RabbitServiceAutoConfiguration::class,
        RabbitMessageChannelBinderConfiguration::class,
        PropertyPlaceholderAutoConfiguration::class,
    )
    @EnableConfigurationProperties(
        RabbitProperties::class,
        RabbitBinderConfigurationProperties::class,
        RabbitExtendedBindingProperties::class
    )
    class RabbitConfig {
    
        @Bean
        @Primary
        @Role(BeanDefinition.ROLE_INFRASTRUCTURE)
        @Order(Ordered.HIGHEST_PRECEDENCE)
        fun customRabbitMessageChannelBinder(
            appCtx: ConfigurableApplicationContext,
            // 其他构造所需的依赖注入
        ): RabbitMessageChannelBinder {
    
            // 移除容器中原生自动配置的binder Bean
            val registry = appCtx.autowireCapableBeanFactory as BeanDefinitionRegistry
            registry.removeBeanDefinition("rabbitMessageChannelBinder")
    
            // 返回自定义Binder实例,初始化逻辑和原生Bean完全一致,仅类为自定义的修改版
            return RabbitMessageChannelBinderWithNoStacktraceRepublished(...)
        }
    }
    

    该方案存在两个互斥的问题:

    • 如果移除原有自动配置的binder Bean,自定义Binder无法读取YAML中配置的参数,比如配置的RabbitMQ地址是localhost:55015,实际会使用默认值localhost:5672,抛出连接拒绝错误,错误日志如下:
      INFO  o.s.a.r.c.CachingConnectionFactory - Attempting to connect to: [localhost:5672]
      INFO  o.s.a.r.l.SimpleMessageListenerContainer - Broker not available; cannot force queue declarations during start: java.net.ConnectException: Connection refused
      
    • 如果不移除原有binder Bean,会抛出多binder冲突错误,错误日志如下:
      Caused by: java.lang.IllegalStateException: Multiple binders are available, however neither default nor per-destination binder name is provided. Available binders are [rabbitMessageChannelBinder, customRabbitMessageChannelBinder]
          at org.springframework.cloud.stream.binder.DefaultBinderFactory.getBinder(DefaultBinderFactory.java:145)
      

运行环境:Spring Cloud Stream 3.1.6、Spring Boot 2.6.6


可行解决方案

不需要重写整个Binder,直接替换Binder内部用于DLQ重发的RepublishMessageRecoverer即可,步骤如下:

  1. 自定义一个继承RepublishMessageRecoverer的恢复器,重写输出消息头的逻辑,在重发前移除x-exception-stacktrace头
  2. 通过RabbitListenerContainerCustomizer获取到每个消费者监听器容器配置的错误处理逻辑,替换其中的恢复器实例为自定义版本

实现代码(Kotlin)

import org.springframework.amqp.core.Message
import org.springframework.amqp.rabbit.retry.RepublishMessageRecoverer
import org.springframework.amqp.rabbit.config.RabbitListenerContainerCustomizer
import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer
import org.springframework.boot.autoconfigure.amqp.RabbitProperties
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Configuration

@Configuration
class DlqHeaderConfig {

    /**
     * 自定义重发恢复器,移除栈追踪头
     */
    class NoStacktraceRepublishRecoverer(
        template: org.springframework.amqp.rabbit.core.RabbitTemplate,
        errorExchange: String?,
        errorRoutingKey: String?
    ) : RepublishMessageRecoverer(template, errorExchange, errorRoutingKey) {
        override fun additionalHeaders(message: Message?, cause: Throwable?) {
            super.additionalHeaders(message, cause)
            // 移除栈追踪请求头
            message?.messageProperties?.removeHeader("x-exception-stacktrace")
        }
    }

    @Bean
    fun rabbitContainerCustomizer(
        rabbitTemplate: org.springframework.amqp.rabbit.core.RabbitTemplate
    ): RabbitListenerContainerCustomizer<SimpleMessageListenerContainer> {
        return RabbitListenerContainerCustomizer { container, _ ->
            val adviceChain = container.adviceChain
            if (adviceChain != null) {
                adviceChain.forEach { advice ->
                    // 找到Binder内置的重试拦截器
                    if (advice is org.springframework.retry.interceptor.RetryOperationsInterceptor) {
                        val recoverer = advice.recoveryCallback
                        if (recoverer is org.springframework.amqp.rabbit.retry.MessageRecoverer) {
                            // 读取原有恢复器的DLQ配置,无需重复硬编码
                            if (recoverer is RepublishMessageRecoverer) {
                                val errorExchange = recoverer.errorExchangeName
                                val errorRoutingKey = recoverer.errorRoutingKey
                                advice.recoveryCallback = org.springframework.retry.interceptor.RetryOperationsInterceptor
                                    .getDefaultRecoveryCallback(
                                        NoStacktraceRepublishRecoverer(rabbitTemplate, errorExchange, errorRoutingKey)
                                    )
                            }
                        }
                    }
                }
            }
        }
    }
}

方案说明

  • 不需要修改原有YAML配置,不需要替换整个Binder实例,不会出现配置读取失败、多Bean冲突的问题
  • 自定义恢复器完全复用原有重发逻辑,仅在添加完默认头后移除不需要的栈追踪头,不会影响其他DLQ功能(比如TTL、死信路由配置)
  • 完全适配Spring Cloud Stream 3.1.x版本,不需要升级依赖

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 11:36:19