生产者与消费者微服务间WebSocket消息传递失败求助
生产者与消费者微服务消息传递失败排查求助
问题现象
两个微服务已正常启动,完成密码验证后,生产者发送的消息无法到达消费者,消息数据未写入数据库。
相关代码
生产者Controller
@RestController @RequestMapping("/api/producer") public class ProducerController { @Autowired private SimpMessagingTemplate messagingTemplate; private UserAuditRepository userAuditRepository; private static final Logger logger = LoggerFactory.getLogger(ProducerController.class); @Autowired public ProducerController(UserAuditRepository userAuditRepository) { this.userAuditRepository = userAuditRepository; } @PostMapping("/send-message/{userId}") public ModelAndView sendMessageToUser(@PathVariable String userId) { logger.info("Sending message to user: {}", userId); // 检查用户是否存在 Optional<UserAudit> existingUser = userAuditRepository.findByUserId(userId); UserAudit userAudit = existingUser.orElseGet(() -> { UserAudit newUser = new UserAudit(); newUser.setUserId(userId); return newUser; }); MessageDto messageDto = new MessageDto(true, "gold", 123, "Hello World", userId); messagingTemplate.convertAndSendToUser(userId, "/topic/messages", messageDto); // 发送消息后重定向到用户审计页面 return new ModelAndView(new RedirectView("/consumersocket/api/consumer/user-audit", true)); } }
消费者Controller
@RestController @RequestMapping("/api/consumer") public class ConsumerController { @Autowired private UserAuditRepository userAuditRepository; @MessageMapping("/messages") @SendToUser("/topic/messages") public void processMessageFromProducer(@Payload MessageDto messageDto, Principal principal) { UserAudit userAudit = new UserAudit(); userAudit.setUserId(messageDto.getUserId()); // 从消息中提取userId userAudit.setIsActive(messageDto.isActive()); userAudit.setColor(messageDto.getColor()); userAudit.setNumber(messageDto.getNumber()); userAudit.setMessage(messageDto.getMessage()); userAuditRepository.save(userAudit); } @GetMapping("/user-audit") public String getUserAudit(Principal principal, Model model) { Optional<UserAudit> userAuditOptional = userAuditRepository.findByUserId(principal.getName()); if (userAuditOptional.isPresent()) { UserAudit userAudit = userAuditOptional.get(); model.addAttribute("userAuditList", Collections.singletonList(userAudit)); } else { model.addAttribute("userAuditList", Collections.emptyList()); } return "user-audit"; } }
服务配置与接口信息
- 生产者服务上下文路径:
server.servlet.context-path=/producersocket - 生产者发送接口:
http://localhost:8080/producersocket/api/producer/send-message/{userId} - 消费者查询接口:
http://localhost:8081/consumersocket/api/consumer/user-audit
已尝试调整接口URL,但问题仍未解决,寻求排查思路。
排查建议
WebSocket核心配置检查
- 确认两个服务都添加了
@EnableWebSocketMessageBroker注解,消息代理(SimpleBroker/StompBrokerRelay)的前缀配置一致,比如生产者和消费者是否都声明了相同的/topic前缀。 - 验证
convertAndSendToUser的路径格式:部分配置下用户专属路径应为/user/{userId}/topic/messages,需确认消息代理是否自动处理前缀拼接。
- 确认两个服务都添加了
跨服务消息代理验证
- 如果使用外部消息代理(RabbitMQ/ActiveMQ),检查两个服务的代理连接参数(地址、虚拟主机、账号密码)是否完全一致。
- 通过代理管理界面查看对应主题/队列的消息流转,确认生产者是否成功投递消息,是否有未被消费的消息积压。
用户订阅与身份校验
- 确认消费者服务或前端已成功订阅目标主题
/user/{userId}/topic/messages,未订阅则无法接收消息。 - 检查消费者
processMessageFromProducer方法中的Principal是否能正确获取用户身份,若用户未通过WebSocket认证,消息路由会失败。
- 确认消费者服务或前端已成功订阅目标主题
日志与调试
- 在生产者的
sendMessageToUser方法末尾添加日志,确认messagingTemplate调用未抛出异常。 - 在消费者的
processMessageFromProducer方法开头添加日志,验证是否进入消息处理逻辑。 - 查看两个服务的日志,排查WebSocket连接失败、代理通信异常等错误信息。
- 在生产者的
网络与跨域检查
- 确认生产者、消费者、消息代理三者之间网络连通,防火墙未拦截对应端口。
- 若涉及前端跨域,检查WebSocket的跨域配置是否允许对应域名的请求。
内容的提问来源于stack exchange,提问作者PoToM
相关产品推荐
相关产品推荐

