404 Not Found

404 Not Found


nginx

WebSocket — リアルタイム双方向通信

HTTPは手紙を送るようなものです。一方通行で戻ってきます。WebSocketは電話をかけるようなものです。両者がいつでも話すことができ, かけ直す必要がありません。

1. 学ぶ内容


2. Aliceのリアルストーリー

(1) ペインポイント:価格変更はポーリングでしか取得できない

Bobのフロントエンドは5秒ごとにPriceTracker APIをポーリングして価格変更を確認していますが, 百万件の製品の99%は5秒以内に変更されないため, 99%のリクエストが無駄になっています。さらに悪いことに, 価格変更が表示されるまで最大5秒かかり, 顧客から「価格がリアルタイムではない」と苦情が出ています。

(2) WebSocketソリューション

WebSocketは永続的な双方向接続を確立し, 価格変更が発生したときにサーバーからフロントエンドにアクティブにプッシュします。ポーリングが不要になり, ゼロ無駄, ゼロレイテンシを実現します。

PYTHON
@app.websocket("/ws/prices")
async def price_websocket(websocket: WebSocket):
    await websocket.accept()
    while True:
        data = await websocket.receive_text()
        await websocket.send_json({"price_update": data})

(3) 収益

ポーリングリクエスト数が1秒200回から0に減少し, 価格更新の遅延が5秒から50msに短縮されました。BobのフロントエンドはAPIクォータを無駄に消費しなくなり, 顧客満足度が大幅に向上しました。


3. WebSocketの基礎

(1) ライフサイクルステート

100%
stateDiagram-v2
    [*] --> CONNECTING: クライアントが開始
    CONNECTING --> CONNECTED: accept()
    CONNECTED --> RECEIVING: receive()
    RECEIVING --> CONNECTED: send()
    CONNECTED --> CLOSING: close() / disconnect
    CLOSING --> CLOSED: 接続閉じる
    CLOSED --> [*]

(1) ▶サンプル:最小のWebSocketエンドポイント

PYTHON
from fastapi import FastAPI, WebSocket

app = FastAPI()

@app.websocket("/ws/echo")
async def websocket_echo(websocket: WebSocket):
    await websocket.accept()  # 接続を受け入れる
    try:
        while True:
            data = await websocket.receive_text()
            await websocket.send_text(f"Echo: {data}")
    except Exception:
        await websocket.close()

出力:

TEXT
# 関数定義成功

(2) ▶サンプル:パスパラメータ付きWebSocket

PYTHON
@app.websocket("/ws/products/{product_id}/prices")
async def product_price_stream(
    websocket: WebSocket,
    product_id: int,
):
    await websocket.accept()
    try:
        while True:
            data = await websocket.receive_json()
            # プロダクトコンテキスト付きでエコーバック
            await websocket.send_json({
                "product_id": product_id,
                "price": data.get("price"),
                "currency": data.get("currency", "USD"),
            })
    except Exception:
        await websocket.close()

出力:

TEXT
# 関数定義成功

4. ConnectionManagerパターン

(1) ConnectionManagerの設計

(1) ▶サンプル:ConnectionManagerの実装

PYTHON
from fastapi import FastAPI, WebSocket
from typing import Dict, List
import json

class ConnectionManager:
    def __init__(self):
        # マッピング:product_id -> アクティブな接続のリスト
        self.active_connections: Dict[int, List[WebSocket]] = {}

    async def connect(self, websocket: WebSocket, product_id: int):
        await websocket.accept()
        if product_id not in self.active_connections:
            self.active_connections[product_id] = []
        self.active_connections[product_id].append(websocket)

    def disconnect(self, websocket: WebSocket, product_id: int):
        if product_id in self.active_connections:
            self.active_connections[product_id].remove(websocket)
            if not self.active_connections[product_id]:
                del self.active_connections[product_id]

    async def broadcast_to_product(self, product_id: int, message: dict):
        if product_id in self.active_connections:
            dead_connections = []
            for connection in self.active_connections[product_id]:
                try:
                    await connection.send_json(message)
                except Exception:
                    dead_connections.append(connection)
            # 切断された接続をクリーンアップ
            for conn in dead_connections:
                self.disconnect(conn, product_id)

    async def broadcast_all(self, message: dict):
        for product_id in list(self.active_connections.keys()):
            await self.broadcast_to_product(product_id, message)

manager = ConnectionManager()

出力:

TEXT
# 関数定義成功

(2) WebSocketプッシュアーキテクチャ

100%
flowchart TD
    Update[Price Update via HTTP] --> Handler[API Handler]
    Handler --> Manager[ConnectionManager]
    Manager --> WS1[WebSocket Client 1]
    Manager --> WS2[WebSocket Client 2]
    Manager --> WSN[WebSocket Client N]
    
    subgraph Subscribers
        WS1
        WS2
        WSN
    end

(2) ▶サンプル:ConnectionManagerとWebSocketエンドポイントの使用

PYTHON
app = FastAPI()

@app.websocket("/ws/products/{product_id}/prices")
async def price_websocket(websocket: WebSocket, product_id: int):
    await manager.connect(websocket, product_id)
    try:
        while True:
            # 接続を維持, クライアントメッセージを受信
            data = await websocket.receive_text()
    except Exception:
        manager.disconnect(websocket, product_id)

出力:

TEXT
# 関数定義成功

5. 認証付きWebSocket

(1) ハンドシェイクフェーズでのJWT検証

WebSocketには標準的なヘッダー機構がありません。トークンはクエリパラメータで渡されます。

(1) ▶サンプル:WebSocket JWT認証

PYTHON
from fastapi import WebSocket, Query, HTTPException
from jose import jwt, JWTError
from app.core.config import settings

async def verify_ws_token(token: str) -> dict:
    try:
        payload = jwt.decode(token, settings.secret_key, algorithms=[settings.algorithm])
        return payload
    except JWTError:
        raise ValueError("Invalid token")

@app.websocket("/ws/prices")
async def authenticated_price_ws(
    websocket: WebSocket,
    token: str = Query(..., description="JWT access token"),
):
    # 接続受け入れ前にトークンを検証
    try:
        user = await verify_ws_token(token)
    except ValueError:
        await websocket.close(code=4001, reason="Authentication failed")
        return

    await websocket.accept()
    try:
        while True:
            data = await websocket.receive_text()
            await websocket.send_json({"user": user.get("sub"), "data": data})
    except Exception:
        pass

出力:

TEXT
# 関数定義成功
認証方式 実装 利点 欠点
クエリパラメータ ?token=xxx シンプル URLログにトークンが残る
最初のメッセージ 接続後にトークンを送信 URLに露出しない 追加のラウンドトリップ
Sec-WebSocket-Protocol サブプロトコル経由でトークン送信 URLに露出しない 非標準的な使用方法

6. HTTPとWebSocketの連携

(1) 価格変更がプッシュ通知をトリガー

(1) ▶サンプル:HTTPエンドポイントがWebSocketブロードキャストをトリガー

PYTHON
from pydantic import BaseModel, Field
from datetime import datetime

class PriceUpdate(BaseModel):
    product_id: int = Field(gt=0)
    price: float = Field(gt=0, description="New price in USD")
    currency: str = Field(default="USD")
    source: str = Field(max_length=100)

app = FastAPI()
manager = ConnectionManager()

@app.post("/api/v1/prices", response_model=PriceResponse)
async def create_price(
    price: PriceCreate,
    db: AsyncSession = Depends(get_db),
    user=Depends(get_current_user),
):
    service = PriceService(db)
    result = await service.create_price(price)
    
    # 価格作成成功後にWebSocketブロードキャストをトリガー
    await manager.broadcast_to_product(
        product_id=price.product_id,
        message={
            "event": "price_update",
            "product_id": price.product_id,
            "new_price": price.price,
            "currency": price.currency,
            "source": price.source,
            "timestamp": datetime.utcnow().isoformat(),
        },
    )
    return result

出力:

TEXT
# 関数定義成功

(2) ▶サンプル:グローバル価格フィード (すべての変更を購読)

PYTHON
@app.websocket("/ws/prices/stream")
async def global_price_stream(websocket: WebSocket):
    await websocket.accept()
    # グローバル購読者リストに追加
    manager.global_connections.append(websocket)
    try:
        while True:
            # クライアントからのハートビート/_pingを受信
            data = await websocket.receive_text()
            if data == "ping":
                await websocket.send_text("pong")
    except Exception:
        manager.global_connections.remove(websocket)

出力:

TEXT
# 関数定義成功

❓ よくある質問

Q WebSocketとSSE (Server-Sent Events)はどう使い分けるべきですか?
A 双方向通信が必要な場合 (チャット, リアルタイムコラボレーションなど)はWebSocketを使用し, サーバー側プッシュのみが必要な場合 (ログストリーム, 通知など)はSSEを使用します。SSEの方がシンプルで自動再接続機能もあります。
Q WebSocket接続数に上限はありますか?
A 単一マシンではファイル記述子の数で制限されます (通常65,535)。本番環境ではNginxでロードバランシングを行い, 複数インスタンスが接続状態を共有します (Redis Pub/Sub経由)。
Q WebSocketの切断はどのように処理すべきですか?
A クライアントは自動再接続 (指数バックオフ)を実装します。サーバーはping/pongハートビートで死活接続を検出し, ConnectionManagerが自動的にクリーンアップします。
Q WebSocketでPydanticを使ってメッセージを検証できますか?
A はい。メッセージ受信後, MessageModel.model_validate(data)で検証します。検証失敗時はエラーメッセージをクライアントに送信します。
Q 複数のワーカープロセスでデータをブロードキャストするにはどうすればよいですか?
A 単一プロセスのインメモリConnectionManagerでは不十分です。本番環境ではRedis Pub/Subを使用します。価格変更をRedisチャネルにパブリッシュし, 各ワーカーがサブスクライブしてローカル接続にプッシュします。
Q WebSocketのCORSはどのように処理しますか?
A ブラウザのWebSocketも同一生成元ポリシーの対象です。FastAPIのCORSMiddlewareがWebSocketのクロスオリジンハンドシェイクを自動的に処理します。

📖 まとめ


📝 練習問題

  1. 基本問題 (難易度 ⭐):WebSocketエコーエンドポイントを作成し, クライアントがテキストを送信するとサーバーがそのまま返します。ブラウザの開発者ツールまたはwscatでテストします。ヒント:@app.websocket("/ws/echo") + accept() + receive_text()
  2. 応用問題 (難易度 ⭐⭐):複数クライアントが同じプロダクトの価格変更を購読できるConnectionManagerを実装します。新しい価格が作成されたとき, そのプロダクトを購読しているすべてのWebSocketクライアントにブロードキャストします。ヒント:Dict[int, List[WebSocket]] + broadcast_to_product()
  3. チャレンジ (難易度 ⭐⭐⭐):WebSocketにJWT認証を追加 (クエリパラメータでトークンを渡す)し, HTTP価格作成エンドポイントでWebSocketブロードキャストをトリガーして, 完全な「HTTP書き込み → WebSocketプッシュ」ワークフローを実装します。ヒント:token: str = Query(...) + verify_ws_token() + HTTPエンドポイントでmanager.broadcast_to_product()を呼び出し

---|

Web-Tutorial.com

Web-Tutorial 技術チーム

複数の開発者によって共同維持されているプログラミングチュートリアルプラットフォーム。各チュートリアルは専門分野の開発者が執筆・レビューしています。正確で信頼性の高いコンテンツを目指しています — 問題を見つけた場合はお知らせください。

100%