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

如何使用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();
    }
  }
}

关键说明

  1. STOMP命令帧:所有交互必须通过STOMP规定的命令(CONNECT/SUBSCRIBE/SEND/DISCONNECT)实现,RxJS不会自动封装这些格式。
  2. 订阅ID:每个SUBSCRIBE命令需要唯一的id头,用于服务器识别订阅关系。
  3. 消息解析:服务器推送的消息是MESSAGE帧,需解析body字段获取实际业务数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 01:05:15