Chuyển đến nội dung chính

レッスン 14: WebSocket とリアルタイム通信

FastAPI の WebSocket エンドポイント、接続マネージャー、ブロードキャスト、ルーム パターン。リアルタイムチャット、通知。 Server-Sent Events (SSE) とロングポーリングの代替手段。

💻 プログラミング — レッスン 14 レッスン 14: WebSocket とリアルタイム コミュニケーション

Python FastAPI: 基本から高度まで

パート 4: 高度な機能

xdev.asia

1. FastAPI の WebSocket の基本

from fastapi import FastAPI, WebSocket, WebSocketDisconnect

app = FastAPI()


@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket):
    """WebSocket endpoint cơ bản."""
    await websocket.accept()
    try:
        while True:
            # Nhận message từ client
            data = await websocket.receive_text()
            # Gửi message lại client
            await websocket.send_text(f"Echo: {data}")
    except WebSocketDisconnect:
        print("Client disconnected")

2. 接続マネージャー

# app/core/websocket.py
import json
from dataclasses import dataclass, field
from fastapi import WebSocket


@dataclass
class ConnectionManager:
    """Quản lý WebSocket connections."""

    # Tất cả connections
    active_connections: list[WebSocket] = field(default_factory=list)

    # Connections theo room
    rooms: dict[str, list[WebSocket]] = field(default_factory=dict)

    # Connections theo user
    user_connections: dict[int, list[WebSocket]] = field(default_factory=dict)

    async def connect(
        self, websocket: WebSocket, user_id: int | None = None, room: str | None = None
    ):
        """Kết nối client."""
        await websocket.accept()
        self.active_connections.append(websocket)

        if user_id:
            self.user_connections.setdefault(user_id, []).append(websocket)

        if room:
            self.rooms.setdefault(room, []).append(websocket)

    def disconnect(
        self, websocket: WebSocket, user_id: int | None = None, room: str | None = None
    ):
        """Ngắt kết nối client."""
        if websocket in self.active_connections:
            self.active_connections.remove(websocket)

        if user_id and user_id in self.user_connections:
            conns = self.user_connections[user_id]
            if websocket in conns:
                conns.remove(websocket)
            if not conns:
                del self.user_connections[user_id]

        if room and room in self.rooms:
            conns = self.rooms[room]
            if websocket in conns:
                conns.remove(websocket)
            if not conns:
                del self.rooms[room]

    async def send_personal(self, message: dict, websocket: WebSocket):
        """Gửi message đến 1 client."""
        await websocket.send_json(message)

    async def send_to_user(self, message: dict, user_id: int):
        """Gửi message đến tất cả connections của 1 user."""
        connections = self.user_connections.get(user_id, [])
        for ws in connections:
            try:
                await ws.send_json(message)
            except Exception:
                pass

    async def broadcast(self, message: dict, exclude: WebSocket | None = None):
        """Gửi message đến tất cả clients."""
        for ws in self.active_connections:
            if ws != exclude:
                try:
                    await ws.send_json(message)
                except Exception:
                    pass

    async def broadcast_to_room(
        self, message: dict, room: str, exclude: WebSocket | None = None
    ):
        """Gửi message đến tất cả clients trong room."""
        connections = self.rooms.get(room, [])
        for ws in connections:
            if ws != exclude:
                try:
                    await ws.send_json(message)
                except Exception:
                    pass

    @property
    def connection_count(self) -> int:
        return len(self.active_connections)


# Global instance
manager = ConnectionManager()

3. リアルタイムチャットルーム

# app/api/v1/chat.py
from fastapi import APIRouter, WebSocket, WebSocketDisconnect, Depends, Query
from datetime import datetime

from app.core.websocket import manager

router = APIRouter(tags=["Chat"])


@router.websocket("/ws/chat/{room_id}")
async def chat_websocket(
    websocket: WebSocket,
    room_id: str,
    username: str = Query(...),
):
    """WebSocket endpoint cho chat room."""
    await manager.connect(websocket, room=room_id)

    # Thông báo user joined
    await manager.broadcast_to_room(
        {
            "type": "system",
            "message": f"{username} joined the room",
            "timestamp": datetime.utcnow().isoformat(),
        },
        room=room_id,
        exclude=websocket,
    )

    try:
        while True:
            # Nhận message
            data = await websocket.receive_json()

            # Broadcast to room
            await manager.broadcast_to_room(
                {
                    "type": "message",
                    "username": username,
                    "message": data.get("message", ""),
                    "timestamp": datetime.utcnow().isoformat(),
                },
                room=room_id,
            )
    except WebSocketDisconnect:
        manager.disconnect(websocket, room=room_id)
        await manager.broadcast_to_room(
            {
                "type": "system",
                "message": f"{username} left the room",
                "timestamp": datetime.utcnow().isoformat(),
            },
            room=room_id,
        )

4. リアルタイム通知

# app/api/v1/notifications.py
from fastapi import APIRouter, WebSocket, WebSocketDisconnect, Depends

from app.core.auth import get_current_user_ws
from app.core.websocket import manager
from app.models.user import User

router = APIRouter(tags=["Notifications"])


@router.websocket("/ws/notifications")
async def notifications_websocket(
    websocket: WebSocket,
    token: str,  # JWT token qua query parameter
):
    """WebSocket cho real-time notifications."""
    # Authenticate WebSocket connection
    from app.core.security import decode_token
    payload = decode_token(token)
    if not payload:
        await websocket.close(code=4001, reason="Unauthorized")
        return

    user_id = int(payload["sub"])
    await manager.connect(websocket, user_id=user_id)

    try:
        while True:
            # Keep connection alive, nhận ping/pong
            data = await websocket.receive_json()
            if data.get("type") == "ping":
                await websocket.send_json({"type": "pong"})
    except WebSocketDisconnect:
        manager.disconnect(websocket, user_id=user_id)


# Service gửi notification
async def send_notification(user_id: int, notification: dict):
    """Gửi notification đến user qua WebSocket."""
    await manager.send_to_user(
        {
            "type": "notification",
            **notification,
        },
        user_id=user_id,
    )

5. サーバー送信イベント (SSE)

# app/api/v1/events.py
import asyncio
from datetime import datetime

from fastapi import APIRouter, Request
from fastapi.responses import StreamingResponse

router = APIRouter(tags=["Events"])


async def event_generator(request: Request):
    """Generate Server-Sent Events."""
    while True:
        # Check nếu client disconnect
        if await request.is_disconnected():
            break

        # Gửi event
        data = {
            "timestamp": datetime.utcnow().isoformat(),
            "connections": 42,  # Example data
        }
        yield f"data: {data}\n\n"

        await asyncio.sleep(1)  # Gửi mỗi giây


@router.get("/events/stream")
async def stream_events(request: Request):
    """SSE endpoint - alternative cho WebSocket."""
    return StreamingResponse(
        event_generator(request),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "Connection": "keep-alive",
            "X-Accel-Buffering": "no",  # Disable nginx buffering
        },
    )

6. WebSocket、SSE、ロングポーリング

特長WebソケットSSEロングポーリング
方向双方向サーバー → クライアントリクエストとレスポンス
プロトコルWS/WSSHTTPHTTP
再接続マニュアル自動(ブラウザ)マニュアル
バイナリデータはいいいえ (テキストのみ)はい
プロキシのサポート時々問題が発生するフルフル
ユースケースチャット、ゲーム通知、フィード簡単なアップデート

概要

この記事では、以下を実装しました。

  • WebSocketエンドポイント:双方向リアルタイム通信
  • 接続マネージャー: 接続、ルーム、ユーザーを管理
  • チャットルーム: ルームパターンによるリアルタイムチャット
  • 通知:WebSocket経由のプッシュ通知
  • SSE: 一方向ストリーミング用のサーバー送信イベント

次の記事では、バックグラウンド タスクとセロリについて学びます。