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

如何实现FastAPI Websocket向多Angular前端同步跨客户端变更?

问题描述

我参考教程实现了Python 3.8 FastAPI后端与Angular的WebSocket连接,单客户端运行正常。但希望实现多个客户端能感知其他客户端触发的后端变更——当前每个客户端会创建独立WebSocket连接,仅能接收自身操作的反馈,无法感知其他客户端的变更。请问如何实现让所有连接的前端客户端都能检测到其他客户端触发的变更?

原代码

FastAPI后端(main.py)

import asyncio
import logging
from datetime import datetime

from fastapi import FastAPI, WebSocket, WebSocketDisconnect

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("FastAPI app")

app = FastAPI()


async def heavy_data_processing(data: dict):
    """Some (fake) heavy data processing logic."""
    await asyncio.sleep(2)
    message_processed = data.get("message", "").upper()
    return message_processed


@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket):
    # Accept the connection from a client.
    await websocket.accept()

    while True:
        try:
            # Receive the JSON data sent by a client.
            data = await websocket.receive_json()
            # Some (fake) heavey data processing logic.
            message_processed = await heavy_data_processing(data)
            # Send JSON data to the client.
            await websocket.send_json(
                {
                    "message": message_processed,
                    "time": datetime.now().strftime("%H:%M:%S"),
                }
            )
        except WebSocketDisconnect:
            logger.info("The connection is closed.")
            break 

Angular前端

app.component.html
<h2>Send a message to the server:</h2>

<form (ngSubmit)="sendMessage(message); message = ''">
    <input [(ngModel)]="message" name="message" type="text" autocomplete="off" />
    <button type="submit" style="margin-left: 10px;">Send</button>
</form>

<h2>Received messages from the server:</h2>
<ul>
    <li *ngFor="let data of webSocketService.receivedData">
        {{ data.time }}: {{ data.message }}
    </li>
</ul>
app.component.ts
import { Component, OnDestroy } from '@angular/core';
import { WebSocketService } from './websocket.service';

@Component({
    selector: 'app-root',
    templateUrl: './app.component.html',
    styleUrls: ['./app.component.css'],
})
export class AppComponent implements OnDestroy {
    message = '';

    constructor(public webSocketService: WebSocketService) {
        this.webSocketService.connect();
    }

    sendMessage(message: string) {
        this.webSocketService.sendMessage(message);
    }

    ngOnDestroy() {
        this.webSocketService.close();
    }
}
websocket.service.ts
import { Injectable } from '@angular/core';
import { webSocket, WebSocketSubject } from 'rxjs/webSocket';
import { environment } from '../environments/environment';

interface MessageData {
    message: string;
    time?: string;
}

@Injectable({
    providedIn: 'root',
})
export class WebSocketService {
    private socket$!: WebSocketSubject<any>;
    public receivedData: MessageData[] = [];

    public connect(): void {
        if (!this.socket$ || this.socket$.closed) {
            this.socket$ = webSocket(environment.webSocketUrl);
    
            this.socket$.subscribe((data: MessageData) => {
                this.receivedData.push(data);
            });
        }
    }

    sendMessage(message: string) {
        this.socket$.next({ message });
    }

    close() {
        this.socket$.complete();
    }
}
解决方案

核心思路是在后端维护一个活跃WebSocket连接的集合,当任意客户端触发后端变更后,将处理后的消息广播给所有已连接的客户端,而非仅回复发起请求的客户端。

一、修改FastAPI后端(main.py)

添加连接集合与并发锁,确保多客户端连接时的线程安全:

import asyncio
import logging
from datetime import datetime
from fastapi import FastAPI, WebSocket, WebSocketDisconnect

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("FastAPI app")

app = FastAPI()

# 存储所有活跃的WebSocket连接
active_connections: list[WebSocket] = []
# 并发锁,防止多客户端同时修改连接集合
connection_lock = asyncio.Lock()


async def heavy_data_processing(data: dict):
    """模拟耗时数据处理"""
    await asyncio.sleep(2)
    message_processed = data.get("message", "").upper()
    return message_processed


async def broadcast_message(message: dict):
    """向所有活跃客户端广播消息"""
    async with connection_lock:
        # 遍历所有连接,发送消息;移除断开的连接
        for connection in active_connections[:]:
            try:
                await connection.send_json(message)
            except WebSocketDisconnect:
                active_connections.remove(connection)
                logger.info("移除断开的客户端连接")


@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket):
    await websocket.accept()
    # 将新连接加入集合
    async with connection_lock:
        active_connections.append(websocket)
        logger.info(f"新客户端连接,当前活跃连接数:{len(active_connections)}")

    try:
        while True:
            data = await websocket.receive_json()
            message_processed = await heavy_data_processing(data)
            response_message = {
                "message": message_processed,
                "time": datetime.now().strftime("%H:%M:%S"),
                # 可选:添加发送方标识,前端可区分消息来源
                # "sender": f"client_{id(websocket)}"
            }
            # 广播给所有客户端
            await broadcast_message(response_message)
    except WebSocketDisconnect:
        # 连接断开时从集合移除
        async with connection_lock:
            active_connections.remove(websocket)
            logger.info(f"客户端断开连接,当前活跃连接数:{len(active_connections)}")

二、前端代码说明

原Angular代码无需修改,因为前端的WebSocket服务已经在监听服务器推送的消息并更新列表。如果需要区分消息是来自自身还是其他客户端,可以在后端的response_message中添加唯一标识(比如客户端连接的ID),前端接收后在界面上标注即可。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 21:05:06