Spring Integration:PubSub订阅者向外部系统发HTTP请求的问题咨询
我尝试通过Spring Integration的服务激活器调用RestTemplate,以POST、DELETE等HTTP方法向外部系统发送消息,现有代码如下:
现有代码实现
PubSub消费者类 MyConsumer
public class MyConsumer{ @Autowired ExternalService externalService; @Bean public PubSubInboundChannelAdapter messageChannelAdapter(final @Qualifier("myInputChannel") MessageChannel inputChannel, PubSubTemplate pubSubTemplate) { PubSubInboundChannelAdapter adapter = new PubSubInboundChannelAdapter(pubSubTemplate, pubSubSubscriptionName); adapter.setOutputChannel(inputChannel); adapter.setAckMode(AckMode.AUTO_ACK); adapter.setErrorChannelName("pubsubErrors"); return adapter; } @ServiceActivator(inputChannel = "pubsubErrors") public void pubsubErrorHandler(Message<MessagingException> exceptionMessage) { BasicAcknowledgeablePubsubMessage originalMessage = (BasicAcknowledgeablePubsubMessage) exceptionMessage .getPayload().getFailedMessage().getHeaders().get(GcpPubSubHeaders.ORIGINAL_MESSAGE); originalMessage.nack(); } @Bean public MessageChannel myInputChannel() { return new DirectChannel(); } @Bean @ServiceActivator(inputChannel = "myInputChannel") public MessageHandler messageReceiver_AddCustomer() { return message -> { externalService.postDataTOExternalSystems(new String((byte[]) message.getPayload())); }; } @Bean @ServiceActivator(inputChannel = "myInputChannel") public MessageHandler messageReceiver_DeleteCustomer() { return message -> { externalService.deleteCustomer(new String((byte[]) message.getPayload())); BasicAcknowledgeablePubsubMessage originalMessage = message.getHeaders().get(GcpPubSubHeaders.ORIGINAL_MESSAGE, BasicAcknowledgeablePubsubMessage.class); originalMessage.ack(); }; } }
外部服务类 ExternalService
public class ExternalService{ void postDataTOExternalSystems(Object obj){ // RequestEntity object formed with HttpEntity object using obj(in json) and headers restTemplate.exchange("https://externalsystems/",HttpMethod.POST,requestEntity,Object.class); } void deleteDatafromExternalSystems(Object obj){ // RequestEntity object formed with HttpEntity object using obj(in json) and headers restTemplate.exchange("https://externalsystems/",HttpMethod.DELETE,requestEntity,Object.class); } }
注意:代码中存在几处语法错误,需先修正:
ExternalService中deleteDatafromExternalSystems方法名与处理器调用的deleteCustomer不匹配HttpMethod.Detele应为HttpMethod.DELETEmessageReceiver_AddCustomer中new String((byte[]) message.getPayload())缺少闭合括号
当前问题:messageReceiver_AddCustomer和messageReceiver_DeleteCustomer共用同一通道,导致执行添加操作时删除处理器也会被触发;同时外部服务调用失败时,控制台会产生无限失败日志。
问题1:单通道下区分不同消息处理器的方案
方案1:SpEL条件注解(最简方案)
直接在@ServiceActivator上添加condition属性,通过消息头或payload内容判断是否执行当前处理器,无需新增通道:
@Bean @ServiceActivator(inputChannel = "myInputChannel", condition = "#headers['operationType'] == 'ADD'") public MessageHandler messageReceiver_AddCustomer() { return message -> { externalService.postDataTOExternalSystems(new String((byte[]) message.getPayload())); }; } @Bean @ServiceActivator(inputChannel = "myInputChannel", condition = "#headers['operationType'] == 'DELETE'") public MessageHandler messageReceiver_DeleteCustomer() { return message -> { externalService.deleteDatafromExternalSystems(new String((byte[]) message.getPayload())); BasicAcknowledgeablePubsubMessage originalMessage = message.getHeaders().get(GcpPubSubHeaders.ORIGINAL_MESSAGE, BasicAcknowledgeablePubsubMessage.class); originalMessage.ack(); }; }
前提:消息发送时需添加
operationType自定义头,或通过SpEL解析payload中的操作标识(比如condition = "new String((byte[])payload).contains('\"action\":\"DELETE\"')")
方案2:消息头路由器
通过HeaderValueRouter将消息路由到对应子通道,逻辑更清晰:
// 配置路由器 @Bean @Router(inputChannel = "myInputChannel") public HeaderValueRouter operationRouter() { HeaderValueRouter router = new HeaderValueRouter("operationType"); router.setChannelMapping("ADD", "addCustomerChannel"); router.setChannelMapping("DELETE", "deleteCustomerChannel"); router.setDefaultOutputChannel("unknownOperationChannel"); // 处理未知操作 return router; } // 子通道定义 @Bean public MessageChannel addCustomerChannel() { return new DirectChannel(); } @Bean public MessageChannel deleteCustomerChannel() { return new DirectChannel(); } // 处理器绑定到对应子通道 @Bean @ServiceActivator(inputChannel = "addCustomerChannel") public MessageHandler messageReceiver_AddCustomer() { // 原有逻辑 } @Bean @ServiceActivator(inputChannel = "deleteCustomerChannel") public MessageHandler messageReceiver_DeleteCustomer() { // 原有逻辑 }
方案3:自定义Payload路由器
如果消息payload本身包含操作标识(比如JSON中的action字段),可自定义路由器解析内容:
@Bean @Router(inputChannel = "myInputChannel") public MessageRouter payloadRouter() { return message -> { String payload = new String((byte[]) message.getPayload()); // 实际项目建议用Jackson等JSON库解析,此处为简化示例 if (payload.contains("\"action\":\"ADD\"")) { return Collections.singletonList("addCustomerChannel"); } else if (payload.contains("\"action\":\"DELETE\"")) { return Collections.singletonList("deleteCustomerChannel"); } return Collections.singletonList("unknownOperationChannel"); }; }
问题2:解决外部服务调用无限失败日志问题
原因分析
错误处理器中调用originalMessage.nack()会将消息重新放回订阅队列,导致消息被重复消费、重复失败,形成无限循环日志。
解决方案
方案1:添加重试机制(Spring Retry)
对外部服务调用添加重试,避免临时故障导致的重复nack:
- 引入依赖:
<dependency> <groupId>org.springframework.retry</groupId> <artifactId>spring-retry</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-aop</artifactId> </dependency>
- 启动类添加
@EnableRetry注解 - 修改
ExternalService方法:
@Retryable(value = {RestClientException.class}, maxAttempts = 3, backoff = @Backoff(delay = 1000)) void postDataTOExternalSystems(Object obj) { // 原有逻辑 } @Retryable(value = {RestClientException.class}, maxAttempts = 3, backoff = @Backoff(delay = 1000)) void deleteDatafromExternalSystems(Object obj) { // 原有逻辑 } // 重试耗尽后的回调 @Recover void recoverPost(RestClientException e, Object obj) { log.error("POST服务重试失败,payload: {}", obj, e); } @Recover void recoverDelete(RestClientException e, Object obj) { log.error("DELETE服务重试失败,payload: {}", obj, e); }
方案2:配置PubSub死信队列
当消息重试超过指定次数后,自动转发到死信队列,不再重复处理:
@Bean public PubSubInboundChannelAdapter messageChannelAdapter(@Qualifier("myInputChannel") MessageChannel inputChannel, PubSubTemplate pubSubTemplate) { // 创建订阅时配置死信策略 Subscription subscription = Subscription.newBuilder() .setName(pubSubSubscriptionName) .setTopic(pubSubTopicName) .setDeadLetterPolicy(DeadLetterPolicy.newBuilder() .setDeadLetterTopic("your-dead-letter-topic") .setMaxDeliveryAttempts(5) // 最多重试5次 .build()) .build(); pubSubTemplate.createSubscription(subscription); PubSubInboundChannelAdapter adapter = new PubSubInboundChannelAdapter(pubSubTemplate, pubSubSubscriptionName); adapter.setOutputChannel(inputChannel); adapter.setAckMode(AckMode.AUTO_ACK); adapter.setErrorChannelName("pubsubErrors"); return adapter; }
方案3:限制错误处理器的nack次数
在错误处理器中判断消息投递次数,超过阈值后不再nack:
@ServiceActivator(inputChannel = "pubsubErrors") public void pubsubErrorHandler(Message<MessagingException> exceptionMessage) { BasicAcknowledgeablePubsubMessage originalMessage = (BasicAcknowledgeablePubsubMessage) exceptionMessage .getPayload().getFailedMessage().getHeaders().get(GcpPubSubHeaders.ORIGINAL_MESSAGE); // 获取PubSub自带的投递次数属性 Integer deliveryAttempt = Integer.parseInt( originalMessage.getPubsubMessage().getAttributesOrDefault("deliveryAttempt", "1") ); if (deliveryAttempt <= 3) { originalMessage.nack(); // 最多重试3次 } else { log.error("消息重试次数超过阈值,messageId: {}", originalMessage.getPubsubMessage().getMessageId()); originalMessage.ack(); // 或转发到死信队列 } }
内容的提问来源于stack exchange,提问作者learner-newfeatures

