如何在失败消息进入死信主题前添加自定义Header?
解决方法
方案一:改用Consumer<Message<MyEventType>>手动处理失败与DLQ发送
将消费者参数从MyEventType改为Message<MyEventType>,获取完整消息头信息,在处理失败时手动构建带自定义Header的消息并发送到DLQ,同时关闭自动DLQ配置。
代码示例
- 修改消费者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); } }; }
- 修改配置关闭自动DLQ:
spring: cloud: stream: kafka: bindings: myConsumerFunction-in-0: consumer: enable-dlq: false # 关闭自动DLQ,改为手动发送
方案二:自定义DeadLetterPublishingRecoverer处理Header修改
结合Spring Retry与自定义死信发布恢复器,在消息重试失败后发送到DLQ前修改Header,贴合Spring Cloud Stream原生重试机制。
代码示例
- 配置重试模板与自定义恢复器:
@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); }
- 修改配置关联重试与恢复器:
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
相关产品推荐
相关产品推荐

