Spring Boot中基于Redis实现跨服务阻塞式消息获取的问题
基于Redis的跨服务阻塞式通信实现问题排查与修复
需求背景
- Service A(Spring Boot):需暴露
/api-publish和/api-consume两个端点,前者向Redis主题topic1发布消息,后者需阻塞等待Service B处理完成后,从topic2收到通知再返回响应 - Service B(ExpressJS):监听
topic1的消息,处理完成后向topic2发送"已处理完成"的通知
现有代码核心问题
- 端点路径不匹配:当前控制器提供的是
/api/v1/public/test,不符合需求的/api-consume,且缺少/api-publish端点 - RedisTemplate泛型不明确:配置的是通配符泛型
RedisTemplate<?,?>,但注入时用RedisTemplate<String, Object>,会引发类型转换警告 - 消息串扰风险:
RedisMessageSubscriber是单例Bean,所有请求共用同一个BlockingQueue,会导致不同请求的处理结果互相干扰 - 无请求标识关联:未在消息中携带唯一请求ID,Service B无法针对性返回处理结果,订阅者无法区分消息归属
修复后的完整实现
1. Redis配置类(修正泛型与命名)
@Configuration public class RedisConfig { @Value("${spring.data.redis.host}") private String redisHost; @Value("${spring.data.redis.port}") private int redisPort; @Bean public RedisConnectionFactory lettuceConnectionFactory() { return new LettuceConnectionFactory(new RedisStandaloneConfiguration(redisHost, redisPort)); } // 明确泛型为<String, String>,匹配字符串消息场景 @Bean public RedisTemplate<String, String> redisTemplate() { RedisTemplate<String, String> template = new RedisTemplate<>(); template.setConnectionFactory(lettuceConnectionFactory()); template.setHashKeySerializer(new StringRedisSerializer()); template.setKeySerializer(new StringRedisSerializer()); template.setHashValueSerializer(new StringRedisSerializer()); template.setValueSerializer(new StringRedisSerializer()); return template; } @Bean MessageListenerAdapter messageListener() { return new MessageListenerAdapter(new RedisMessageSubscriber()); } @Bean ChannelTopic topic2() { return new ChannelTopic("topic2"); } @Bean RedisMessageListenerContainer redisContainer() { RedisMessageListenerContainer container = new RedisMessageListenerContainer(); container.setConnectionFactory(lettuceConnectionFactory()); container.addMessageListener(messageListener(), topic2()); return container; } }
2. 消息订阅服务(解决并发串扰)
通过ConcurrentHashMap维护请求ID与专属阻塞队列的映射,确保每个请求只接收自己的处理结果:
@Service public class RedisMessageSubscriber implements MessageListener { // 键:请求唯一ID,值:对应请求的阻塞队列 private final ConcurrentHashMap<String, BlockingQueue<String>> requestQueueMap = new ConcurrentHashMap<>(); @Override public void onMessage(Message message, byte[] pattern) { String receivedMsg = new String(message.getBody()); // 按约定格式拆分请求ID和处理结果 String[] msgParts = receivedMsg.split(":", 2); if (msgParts.length == 2) { String requestId = msgParts[0]; String result = msgParts[1]; BlockingQueue<String> targetQueue = requestQueueMap.get(requestId); if (targetQueue != null) { targetQueue.offer(result); requestQueueMap.remove(requestId); // 清理已完成的请求映射,避免内存泄漏 } } } // 根据请求ID等待处理结果 public String waitForResult(String requestId, long timeout, TimeUnit unit) throws InterruptedException { BlockingQueue<String> queue = new LinkedBlockingQueue<>(); requestQueueMap.put(requestId, queue); try { String result = queue.poll(timeout, unit); if (result == null) { requestQueueMap.remove(requestId); // 超时后清理无效映射 } return result; } catch (InterruptedException e) { requestQueueMap.remove(requestId); // 中断后清理映射 Thread.currentThread().interrupt(); throw e; } } // 生成唯一请求ID public String generateRequestId() { return UUID.randomUUID().toString(); } }
3. 控制器(实现需求端点)
@RestController public class MessageController { private final RedisTemplate<String, String> redisTemplate; private final RedisMessageSubscriber subscriber; // 构造注入替代@Autowired,符合Spring最佳实践 public MessageController(RedisTemplate<String, String> redisTemplate, RedisMessageSubscriber subscriber) { this.redisTemplate = redisTemplate; this.subscriber = subscriber; } @PostMapping("/api-publish") public ResponseEntity<String> publishMessage(@RequestBody String content) { String requestId = subscriber.generateRequestId(); // 消息格式:requestId:业务内容,让Service B处理后携带相同ID返回 String message = requestId + ":" + content; redisTemplate.convertAndSend("topic1", message); return ResponseEntity.ok("消息已发布,请求ID:" + requestId); } @GetMapping("/api-consume") public ResponseEntity<String> consumeResult(@RequestParam String requestId) { try { String result = subscriber.waitForResult(requestId, 5, TimeUnit.SECONDS); if (result != null) { return ResponseEntity.ok("收到Service B处理结果:" + result); } else { return ResponseEntity.status(HttpStatus.REQUEST_TIMEOUT).body("等待处理结果超时"); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body("获取处理结果失败"); } } }
4. Service B的ExpressJS核心实现
const redis = require('redis'); // 创建订阅者和发布者客户端 const subscriber = redis.createClient({ host: '你的Redis地址', port: 6379 }); const publisher = redis.createClient({ host: '你的Redis地址', port: 6379 }); subscriber.subscribe('topic1'); // 监听topic1消息并处理 subscriber.on('message', (channel, message) => { const [requestId, content] = message.split(':', 2); // 模拟业务处理逻辑 console.log(`处理业务内容:${content}`); // 处理完成后向topic2发送带请求ID的通知 publisher.publish('topic2', `${requestId}:处理完成,内容:${content}`); });
关键修正总结
- 对齐需求的端点路径,补全
/api-publish实现 - 引入请求ID机制,解决多请求场景下的消息串扰问题
- 明确RedisTemplate泛型,消除类型转换风险
- 优化订阅者的队列管理,避免内存泄漏
内容的提问来源于stack exchange,提问作者ginbarca
相关产品推荐
相关产品推荐

