如何在Cloud Run上用Spring Boot+Google Pub/Sub实现WebSocket房间消息广播?
在Cloud Run上用Spring Boot + Google Pub/Sub实现WebSocket会话广播
核心原理
和Node.js+Redis的方案逻辑一致:
- 跨实例广播:借助Google Pub/Sub的扇出能力,将消息推送到所有运行中的Cloud Run实例
- 本地连接映射:每个服务实例用线程安全的数据结构(如
ConcurrentHashMap)跟踪当前连接的WebSocket会话,收到Pub/Sub消息后,仅转发给本地属于目标会话/房间的连接
分步实现
1. 依赖配置
添加Spring WebSocket和Google Cloud Pub/Sub Starter依赖(以Maven为例):
<dependencies> <!-- Spring WebSocket --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId> </dependency> <!-- Google Cloud Pub/Sub --> <dependency> <groupId>com.google.cloud</groupId> <artifactId>spring-cloud-gcp-starter-pubsub</artifactId> </dependency> </dependencies>
2. WebSocket连接管理
创建全局连接管理器,维护会话ID到WebSocketSession的映射,处理连接生命周期:
@Component public class WebSocketConnectionManager { private final ConcurrentHashMap<String, WebSocketSession> sessionMap = new ConcurrentHashMap<>(); public void addSession(String sessionId, WebSocketSession session) { sessionMap.put(sessionId, session); } public void removeSession(String sessionId) { sessionMap.remove(sessionId); } public Optional<WebSocketSession> getSession(String sessionId) { return Optional.ofNullable(sessionMap.get(sessionId)); } }
3. WebSocket端点配置
定义WebSocket端点,处理连接建立、断开和消息接收:
@Configuration @EnableWebSocket public class WebSocketConfig implements WebSocketConfigurer { private final WebSocketConnectionManager connectionManager; public WebSocketConfig(WebSocketConnectionManager connectionManager) { this.connectionManager = connectionManager; } @Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { registry.addHandler(new ChatWebSocketHandler(), "/ws/chat") .setAllowedOrigins("*"); // 根据实际环境调整跨域规则 } private class ChatWebSocketHandler extends TextWebSocketHandler { @Override public void afterConnectionEstablished(WebSocketSession session) throws Exception { // 从会话参数中获取目标会话ID(比如前端连接时携带) String targetSessionId = session.getUri().getQuery().split("=")[1]; connectionManager.addSession(targetSessionId, session); } @Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { String targetSessionId = session.getUri().getQuery().split("=")[1]; connectionManager.removeSession(targetSessionId); } } }
4. Pub/Sub消息订阅与转发
创建Pub/Sub消息接收器,收到消息后转发给对应的WebSocket会话:
@Component public class PubSubMessageReceiver { private final WebSocketConnectionManager connectionManager; public PubSubMessageReceiver(WebSocketConnectionManager connectionManager) { this.connectionManager = connectionManager; } @PubSubSubscriber(subscription = "chat-message-sub") // 替换为你的Pub/Sub订阅名 public void receiveMessage(String payload) throws IOException { // 解析消息:payload需包含targetSessionId和content字段,比如JSON格式 JSONObject message = new JSONObject(payload); String targetSessionId = message.getString("targetSessionId"); String content = message.getString("content"); // 转发给本地对应的WebSocket会话 connectionManager.getSession(targetSessionId).ifPresent(session -> { if (session.isOpen()) { try { session.sendMessage(new TextMessage(content)); } catch (IOException e) { // 处理发送异常,比如移除失效连接 connectionManager.removeSession(targetSessionId); } } }); } }
5. 消息发布逻辑
创建消息发布器,向Pub/Sub主题发送消息:
@Component public class PubSubMessagePublisher { private final PubSubTemplate pubSubTemplate; public PubSubMessagePublisher(PubSubTemplate pubSubTemplate) { this.pubSubTemplate = pubSubTemplate; } public void sendMessageToSession(String targetSessionId, String content) { JSONObject payload = new JSONObject(); payload.put("targetSessionId", targetSessionId); payload.put("content", content); pubSubTemplate.publish("chat-message-topic", payload.toString()); // 替换为你的Pub/Sub主题名 } }
Cloud Run部署注意事项
- 长连接超时:部署时设置
--timeout=86400s(最大允许值),避免Cloud Run主动断开长连接 - 权限配置:为Cloud Run服务账号授予
roles/pubsub.publisher和roles/pubsub.subscriber权限 - 实例缩容处理:利用Cloud Run的
SIGTERM信号,在实例终止前清理WebSocket连接并更新映射
简化方案的工具选择
无需额外第三方库,直接利用Spring Cloud GCP Pub/Sub和Spring WebSocket即可实现,两者的集成已经做了封装,无需手动处理Pub/Sub的连接和订阅细节。
内容的提问来源于stack exchange,提问作者Paul
相关产品推荐
相关产品推荐

