You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

生产者与消费者微服务间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,但问题仍未解决,寻求排查思路。


排查建议

  1. WebSocket核心配置检查

    • 确认两个服务都添加了@EnableWebSocketMessageBroker注解,消息代理(SimpleBroker/StompBrokerRelay)的前缀配置一致,比如生产者和消费者是否都声明了相同的/topic前缀。
    • 验证convertAndSendToUser的路径格式:部分配置下用户专属路径应为/user/{userId}/topic/messages,需确认消息代理是否自动处理前缀拼接。
  2. 跨服务消息代理验证

    • 如果使用外部消息代理(RabbitMQ/ActiveMQ),检查两个服务的代理连接参数(地址、虚拟主机、账号密码)是否完全一致。
    • 通过代理管理界面查看对应主题/队列的消息流转,确认生产者是否成功投递消息,是否有未被消费的消息积压。
  3. 用户订阅与身份校验

    • 确认消费者服务或前端已成功订阅目标主题/user/{userId}/topic/messages,未订阅则无法接收消息。
    • 检查消费者processMessageFromProducer方法中的Principal是否能正确获取用户身份,若用户未通过WebSocket认证,消息路由会失败。
  4. 日志与调试

    • 在生产者的sendMessageToUser方法末尾添加日志,确认messagingTemplate调用未抛出异常。
    • 在消费者的processMessageFromProducer方法开头添加日志,验证是否进入消息处理逻辑。
    • 查看两个服务的日志,排查WebSocket连接失败、代理通信异常等错误信息。
  5. 网络与跨域检查

    • 确认生产者、消费者、消息代理三者之间网络连通,防火墙未拦截对应端口。
    • 若涉及前端跨域,检查WebSocket的跨域配置是否允许对应域名的请求。

内容的提问来源于stack exchange,提问作者PoToM

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 00:56:20