如何使用RxJS webSocket连接Spring Boot WebSocket主题?
如何用RxJS WebSocket对接Spring Boot的STOMP WebSocket服务器?
问题背景
现有一个基于Spring Boot搭建的STOMP协议WebSocket服务器,使用StompJS可以正常连接、订阅主题并收发消息,但改用RxJS webSocket实现时,仅能建立WebSocket连接,无法正常收发数据。
Spring Boot服务器代码
WebSocketConfig类
@Configuration @EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { @Override public void configureMessageBroker(MessageBrokerRegistry config) { config.enableSimpleBroker("/game-updates"); config.setApplicationDestinationPrefixes("/game"); } @Override public void registerStompEndpoints(StompEndpointRegistry registry) { registry.addEndpoint("/ws").setAllowedOriginPatterns("*"); registry.addEndpoint("/ws").setAllowedOriginPatterns("*").withSockJS(); } }
GameController类
@RestController public class GameController { @MessageMapping("/test") @SendTo("/game-updates/test-res") public Greeting greeting(String msg) throws Exception { System.out.println("msg received: " + msg); String msgg = "Hello there: " + msg; return new Greeting(msgg); } }
可用的StompJS实现代码
const stompClient = new StompJs.Client({ brokerURL: 'ws://localhost:8080/ws' }); stompClient.onConnect = (frame) => { setConnected(true); console.log('Connected: ' + frame); stompClient.subscribe('/game-updates/test-res', (greeting) => { showGreeting(JSON.parse(greeting.body).content); }); }; function connect() { stompClient.activate(); } function sendData() { stompClient.publish({ destination: "/game/test", body: JSON.stringify({ 'data': $("#data").val() }) }); }
问题代码(RxJS实现)
import { webSocket } from 'rxjs/webSocket'; @Component({ selector: 'app-game-screen', standalone: true, templateUrl: './game-screen.component.html', styleUrl: './game-screen.component.css', imports: [MainMenuComponent, OptionsComponent, GameBoardComponent, ChatComponent, BoardUpdatesComponent, MatProgressSpinnerModule] }) export class GameScreenComponent implements OnInit { connectToServer() { console.log("connecting to server"); let url: string = "ws://localhost:8080/ws"; this.webSocket = webSocket(url); this.webSocket.subscribe("/game-updates/test-res", { next: (data: any) => { console.log('data from server received:'); console.log(data); }, error: (err: any) => { console.log(err) }, complete: () => { console.log('complete'); } }); this.webSocket.next({ destination: "/game/test", body: JSON.stringify({ 'data': "Hello server" }) }); } }
解决方案:RxJS适配STOMP协议
RxJS的webSocket默认处理原始WebSocket消息,而Spring Boot服务器使用STOMP协议,需要手动构造符合STOMP格式的命令帧才能正常交互。
完整RxJS实现代码
import { webSocket, WebSocketSubject } from 'rxjs/webSocket'; import { tap } from 'rxjs/operators'; import { OnInit, OnDestroy } from '@angular/core'; @Component({ selector: 'app-game-screen', standalone: true, templateUrl: './game-screen.component.html', styleUrl: './game-screen.component.css', imports: [MainMenuComponent, OptionsComponent, GameBoardComponent, ChatComponent, BoardUpdatesComponent, MatProgressSpinnerModule] }) export class GameScreenComponent implements OnInit, OnDestroy { private webSocket$: WebSocketSubject<any>; ngOnInit(): void { this.connectToServer(); } connectToServer() { console.log("connecting to server"); const url = "ws://localhost:8080/ws"; this.webSocket$ = webSocket(url); // 订阅消息流,处理所有STOMP帧 this.webSocket$.pipe( tap((frame) => this.handleStompFrame(frame)) ).subscribe({ error: (err) => console.log('WebSocket error:', err), complete: () => console.log('WebSocket connection closed') }); // 发送STOMP连接帧 this.webSocket$.next({ command: 'CONNECT', headers: { 'accept-version': '1.1,1.0', 'heart-beat': '10000,10000' } }); } private handleStompFrame(frame: any) { // 处理服务器返回的连接成功响应 if (frame.command === 'CONNECTED') { console.log('STOMP protocol connected:', frame); // 订阅目标主题 this.subscribeToTopic('/game-updates/test-res'); // 发送测试消息 this.sendToServer('/game/test', "Hello server"); } // 处理服务器推送的消息 if (frame.command === 'MESSAGE') { const messageBody = JSON.parse(frame.body); console.log('Received server message:', messageBody.content); // 此处可添加消息业务逻辑 } } private subscribeToTopic(topic: string) { this.webSocket$.next({ command: 'SUBSCRIBE', headers: { id: `sub-${Date.now()}`, // 唯一订阅ID destination: topic } }); } private sendToServer(destination: string, data: string) { this.webSocket$.next({ command: 'SEND', headers: { destination: destination }, body: JSON.stringify({ 'data': data }) }); } ngOnDestroy(): void { if (this.webSocket$) { // 发送STOMP断开连接命令 this.webSocket$.next({ command: 'DISCONNECT' }); this.webSocket$.complete(); } } }
关键说明
- STOMP命令帧:所有交互必须通过STOMP规定的命令(
CONNECT/SUBSCRIBE/SEND/DISCONNECT)实现,RxJS不会自动封装这些格式。 - 订阅ID:每个
SUBSCRIBE命令需要唯一的id头,用于服务器识别订阅关系。 - 消息解析:服务器推送的消息是
MESSAGE帧,需解析body字段获取实际业务数据。
内容的提问来源于stack exchange,提问作者Himanshu Sajwan
相关产品推荐
相关产品推荐

