如何启用死信交换机/主题自动创建(兼容Spring Cloud Stream追踪)
Spring Cloud Stream 函数式Bean配置自定义死信监听器(Rabbit)并兼容链路追踪
一、Rabbit死信队列基础配置
先给原输入绑定补充死信相关配置,让处理失败的消息路由到指定的死信交换器:
# 原业务输入绑定配置(保留原有内容) spring.cloud.stream.function.bindings.myFunctionBean-in-0=input spring.cloud.stream.function.bindings.myFunctionBean-out-0=output spring.cloud.stream.bindings.input.destination=my-listening-topic spring.cloud.stream.bindings.input.group=my-app-identifier spring.cloud.stream.bindings.output.destination=topic-i-publish-to # 新增死信相关配置 spring.cloud.stream.bindings.input.consumer.max-attempts=3 spring.cloud.stream.bindings.input.consumer.back-off-initial-interval=1000 spring.cloud.stream.bindings.input.consumer.back-off-multiplier=2 spring.cloud.stream.rabbit.bindings.input.consumer.dead-letter-exchange=my-dlx-exchange spring.cloud.stream.rabbit.bindings.input.consumer.dead-letter-routing-key=my-dlx-routing-key
二、自定义死信监听器(函数式实现)
通过独立的函数式Bean监听死信队列,同时保证链路上下文的传递:
1. 死信输入绑定配置
# 死信监听器的函数绑定 spring.cloud.stream.function.bindings.deadLetterConsumer-in-0=dlx-input # 死信输入绑定关联死信交换器 spring.cloud.stream.bindings.dlx-input.destination=my-dlx-exchange spring.cloud.stream.bindings.dlx-input.group=my-app-dlx-group # Rabbit专属:绑定死信路由键 spring.cloud.stream.rabbit.bindings.dlx-input.consumer.binding-routing-key=my-dlx-routing-key
2. 死信处理函数Bean实现
import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.function.Consumer; @Configuration public class DeadLetterHandlerConfig { private static final Logger log = LoggerFactory.getLogger(DeadLetterHandlerConfig.class); @Bean public Consumer<String> deadLetterConsumer() { return deadMessage -> { // 自定义死信处理逻辑:日志记录、告警推送、持久化入库等 log.error("处理死信消息,内容摘要: {}", deadMessage.substring(0, 100)); // 如需手动重试,可将消息重新发送回业务队列(需避免循环死信) // messageProducer.send(MessageBuilder.withPayload(deadMessage).build()); }; } }
三、兼容Sleuth/Micrometer链路追踪
Spring Cloud Stream默认支持Sleuth链路上下文传递,只需确保依赖正确并启用相关配置:
1. 核心依赖(Maven)
<dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-sleuth</artifactId> </dependency> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-stream-rabbit</artifactId> </dependency> <dependency> <groupId>io.micrometer</groupId> <artifactId>micrometer-tracing-bridge-brave</artifactId> </dependency>
2. 关键注意事项
- RabbitMQ死信转发会自动保留原消息的
b3/traceparent链路头,无需额外配置即可关联原业务链路 - Micrometer会自动收集消息处理的追踪指标,可通过Actuator端点
/actuator/metrics/spring.cloud.stream.function.invocations查看 - 如需自定义追踪标签,可注入
Tracer对象手动添加:
import io.micrometer.tracing.Tracer; import org.springframework.beans.factory.annotation.Autowired; // 在死信处理逻辑中注入Tracer @Autowired private Tracer tracer; // 添加自定义追踪标签 tracer.currentSpan().tag("dead.message.source", "my-listening-topic");
四、验证方式
- 发送一条会处理失败的消息到
my-listening-topic,观察重试3次后是否进入死信队列 - 查看日志中的追踪ID,确认死信处理链路与原业务链路的关联关系
- 通过监控平台查看Micrometer收集的消息处理指标
内容的提问来源于stack exchange,提问作者Dave Ankin
相关产品推荐
相关产品推荐

