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

Angular+Spring Boot实时聊天消息需手动刷新才显示问题排查

Angular + Spring Boot 实时聊天应用实时消息不显示问题

我用Angular和Spring Boot开发实时聊天应用,消息能正常发送并存入数据库,但客户端必须手动刷新页面才能看到新消息。我已经在ChatService里用SockJS和STOMP建立了WebSocket连接,ChatComponent也订阅了消息,但新消息没法立刻显示在聊天窗口。

以下是我的代码实现:

ChatService 代码

@Injectable({
  providedIn: 'root'
})
export class ChatService {
  private apiUrl = 'http://localhost:8080/api/chat';
  private stompClient: Stomp.Client;
  private messagesSubject = new Subject<ChatMessageDTO>();
  public messages$ = this.messagesSubject.asObservable();

  constructor(private http: HttpClient) {}

  initializeWebSocketConnection() {
    const socket = new SockJS('http://localhost:8080/websockets');
    this.stompClient = Stomp.over(socket);
    this.stompClient.connect({}, frame => {
      console.log('Connected: ' + frame);
      const userMatricule = localStorage.getItem('matricule');
      if (userMatricule) {
        this.stompClient.subscribe('/user/' + userMatricule + '/queue/messages', message => {
          console.log('Received message:', message.body);
          this.messagesSubject.next(JSON.parse(message.body));
        });
      }
    });
  }

  sendMessage(message: ChatMessageDTO): Observable<void> {
    return new Observable<void>(observer => {
      this.stompClient.send('/app/chat', {}, JSON.stringify(message));
      observer.next();
      observer.complete();
    });
  }

  getChatHistory(senderId: number, recipientId: number): Observable<ChatMessageDTO[]> {
    const params = new HttpParams()
      .set('senderId', senderId.toString())
      .set('recipientId', recipientId.toString());
    return this.http.get<ChatMessageDTO[]>(`${this.apiUrl}/history`, { params }).pipe(
      catchError(this.handleError<ChatMessageDTO[]>('getChatHistory', []))
    );
  }

  getUniqueConversations(userMatricule: number): Observable<number[]> {
    const params = new HttpParams().set('userMatricule', userMatricule.toString());
    return this.http.get<number[]>(`${this.apiUrl}/unique-conversations`, { params }).pipe(
      catchError(this.handleError<number[]>('getUniqueConversations', []))
    );
  }

  closeSocket() {
    if (this.stompClient) {
      this.stompClient.disconnect(() => {
        console.log('Disconnected');
      });
    }
  }

  private handleError<T>(operation = 'operation', result?: T) {
    return (error: any): Observable<T> => {
      console.error(`${operation} failed: ${error.message}`);
      return new Observable<T>(subscriber => {
        subscriber.next(result as T);
        subscriber.complete();
      });
    };
  }
}

ChatComponent 代码

export class ChatComponent implements OnInit, OnDestroy {
  messages: ChatMessageDTO[] = [];
  archivedConversations: { userId: number }[] = [];
  newMessage: string = '';
  recipientId: number | null = null;
  userMatricule: number = 0;
  private messagesSubscription: Subscription = new Subscription();

  constructor(private chatService: ChatService, private changeDetector: ChangeDetectorRef) {}

  ngOnInit() {
    this.userMatricule = +localStorage.getItem('matricule')!;
    this.chatService.initializeWebSocketConnection();

    this.messagesSubscription = this.chatService.messages$.subscribe(message => {
      if (message) {
        this.messages.push(message);
        console.log('New message received:', message);
        this.changeDetector.detectChanges();
      }
    });

    this.loadArchivedConversations();
  }

  ngOnDestroy() {
    if (this.messagesSubscription) {
      this.messagesSubscription.unsubscribe();
    }
    this.chatService.closeSocket();
  }

  loadArchivedConversations() {
    this.chatService.getUniqueConversations(this.userMatricule).subscribe(conversations => {
      this.archivedConversations = conversations.map(userId => ({
        userId: userId
      }));
    });
  }

  selectRecipient(recipientId: number) {
    this.recipientId = recipientId;
    this.loadMessages();
  }

  loadMessages() {
    if (this.recipientId !== null) {
      this.chatService.getChatHistory(this.userMatricule, this.recipientId).subscribe(messages => {
        this.messages = messages;
        this.changeDetector.detectChanges();
      });
    }
  }

  sendMessage() {
    if (this.newMessage.trim() && this.recipientId !== null) {
      const message: ChatMessageDTO = new ChatMessageDTO(
        0,
        this.userMatricule,
        this.recipientId,
        this.newMessage,
        new Date(),
        true
      );
      this.chatService.sendMessage(message).subscribe(() => {
        this.newMessage = '';
        this.messages.push(message);
        this.changeDetector.detectChanges();
      });
    }
  }
}

前端模板代码

<div class="chat-container">
  <div class="chat-sidebar">
    <h3>Conversations Archivées</h3>
    <ul>
      <li *ngFor="let conv of archivedConversations" (click)="selectRecipient(conv.userId)">
        Conversation avec {{ conv.userId }}
      </li>
    </ul>
  </div>
  <div class="chat-main">
    <div class="chat-header">
      <input type="text" placeholder="Matricule du destinataire" [(ngModel)]="recipientId" />
      <button (click)="loadMessages()">Charger l'historique</button>
    </div>
    <div class="chat-messages">
      <div *ngFor="let message of messages">
        <div class="chat-message" [ngClass]="{'sent': message.senderId === userMatricule, 'received': message.senderId !== userMatricule}">
          <p>{{ message.content }}</p>
          <span>{{ message.timestamp | date:'short' }}</span>
        </div>
      </div>
    </div>
    <div class="chat-input">
      <input type="text" placeholder="Entrez votre message" [(ngModel)]="newMessage" />
      <button (click)="sendMessage()">Envoyer</button>
    </div>
  </div>
</div>

WebSocketConfig 代码

@Configuration
@EnableWebSocketMessageBroker
public class WebSocketConfig implements WebSocketMessageBrokerConfigurer {

    @Override
    public void configureMessageBroker(MessageBrokerRegistry config) {
        config.enableSimpleBroker("/topic", "/queue");
        config.setApplicationDestinationPrefixes("/app");
    }

    @Override
    public void registerStompEndpoints(StompEndpointRegistry registry) {
        registry.addEndpoint("/websockets").setAllowedOrigins("http://localhost:4200").withSockJS();

    }
}

SecurityFilterChain 代码

@Bean
    public SecurityFilterChain securityFilterChain(HttpSecurity http) throws Exception {
        http
                .csrf(AbstractHttpConfigurer::disable)
                .cors(cors -> cors.configurationSource(corsConfigurationSource()))
                .authorizeHttpRequests(authz -> authz
                        .requestMatchers("/api/auth/**").permitAll()
                        .requestMatchers("/actuator/**").permitAll()
                        .requestMatchers("/websockets/**").permitAll()
                        .requestMatchers("/ws/**").permitAll()
                        .anyRequest().authenticated()
                )
                .sessionManagement(session -> session.sessionCreationPolicy(SessionCreationPolicy.STATELESS))
                .headers(headers -> headers.frameOptions().disable())
                .addFilterBefore(jwtAuthenticationFilter, UsernamePasswordAuthenticationFilter.class);
        return http.build();
    }

    @Bean
    public CorsConfigurationSource corsConfigurationSource() {
        CorsConfiguration corsConfiguration = new CorsConfiguration();
        corsConfiguration.setAllowCredentials(true);
        corsConfiguration.setAllowedOrigins(Arrays.asList("http://localhost:4200"));
        corsConfiguration.setAllowedHeaders(Arrays.asList("Origin", "Access-Control-Allow-Origin", "Content-Type",
                "Accept", "Authorization", "Origin, Accept", "X-Requested-With",
                "Access-Control-Request-Method", "Access-Control-Request-Headers", "Upgrade")); // Ajouter "Upgrade"
        corsConfiguration.setExposedHeaders(Arrays.asList("Origin", "Content-Type", "Accept", "Authorization",
                "Access-Control-Allow-Origin", "Access-Control-Allow-Origin", "Access-Control-Allow-Credentials"));
        corsConfiguration.setAllowedMethods(Arrays.asList("GET", "POST", "PUT", "DELETE", "OPTIONS"));

        UrlBasedCorsConfigurationSource source = new UrlBasedCorsConfigurationSource();
        source.registerCorsConfiguration("/**", corsConfiguration);
        return source;
    }

问题原因及解决方案

原因1:订阅消息未过滤当前对话

你当前将所有收到的消息直接加入messages数组,但组件仅显示选中对话的消息。如果收到的消息属于其他对话,不仅不会显示,还会导致数据混乱;即使是当前对话的消息,也需要验证归属后再添加。

解决方案:修改ChatComponent的订阅逻辑,只添加当前选中对话的消息:

this.messagesSubscription = this.chatService.messages$.subscribe(message => {
  if (message && this.recipientId !== null) {
    // 验证消息是否属于当前对话(双向匹配:发件人/收件人对应当前用户/选中的收件人)
    const isCurrentConv = 
      (message.senderId === this.userMatricule && message.recipientId === this.recipientId) ||
      (message.senderId === this.recipientId && message.recipientId === this.userMatricule);
    if (isCurrentConv) {
      this.messages.push(message);
      console.log('New message received:', message);
      this.changeDetector.detectChanges();
    }
  }
});

原因2:后端未正确推送消息至目标用户队列

检查后端消息处理逻辑,确认收到/app/chat消息后,是否将消息发送到收件人的/user/{matricule}/queue/messages队列,而非仅发送给发件人。

正确后端实现:使用SimpMessagingTemplate手动指定收件人:

@Autowired
private SimpMessagingTemplate messagingTemplate;

@MessageMapping("/chat")
public void handleMessage(@Payload ChatMessageDTO message) {
    // 保存消息到数据库
    chatMessageRepository.save(message);
    // 推送消息给收件人
    messagingTemplate.convertAndSendToUser(
        String.valueOf(message.getRecipientId()), 
        "/queue/messages", 
        message
    );
    // 可选:推送消息给发件人,确保发件人实时看到自己的消息(也可以依赖前端手动添加)
    messagingTemplate.convertAndSendToUser(
        String.valueOf(message.getSenderId()), 
        "/queue/messages", 
        message
    );
}

原因3:WebSocket连接或订阅失败

检查浏览器控制台:

  • 确认localStorage.getItem('matricule')能正确获取用户ID,否则订阅逻辑不会执行
  • 查看是否有WebSocket连接错误、跨域报错(尽管配置了CORS,但需确认Upgrade头是否被允许)

原因4:Subject特性导致消息丢失

Subject不会保留历史消息,如果组件在WebSocket连接建立后才订阅,可能错过早期消息。改用BehaviorSubject可以保留最新消息,确保组件订阅后能立即获取。

修改ChatService:

private messagesSubject = new BehaviorSubject<ChatMessageDTO | null>(null);
public messages$ = this.messagesSubject.asObservable().pipe(filter(message => message !== null));

额外优化:自动更新对话列表

收到新消息时,如果对话不在归档列表中,自动添加进去:

this.messagesSubscription = this.chatService.messages$.subscribe(message => {
  if (message && this.recipientId !== null) {
    // 验证当前对话消息并添加
    const isCurrentConv = 
      (message.senderId === this.userMatricule && message.recipientId === this.recipientId) ||
      (message.senderId === this.recipientId && message.recipientId === this.userMatricule);
    if (isCurrentConv) {
      this.messages.push(message);
    }
    // 更新归档对话列表
    const convId = message.senderId === this.userMatricule ? message.recipientId : message.senderId;
    const exists = this.archivedConversations.some(conv => conv.userId === convId);
    if (!exists) {
      this.archivedConversations.push({ userId: convId });
    }
    this.changeDetector.detectChanges();
  }
});

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 19:57:32