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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 23:07:50