WebSocket — リアルタイム双方向通信
HTTPは手紙を送るようなものです。一方通行で戻ってきます。WebSocketは電話をかけるようなものです。両者がいつでも話すことができ, かけ直す必要がありません。
1. 学ぶ内容
- WebSocketの基礎:
@app.websocketデコレータとライフサイクル (connect/receive/disconnect) - ConnectionManagerパターン:
ConnectionManagerクラスの設計とブロードキャスト機構 - 認証付きWebSocket:ハンドシェイクフェーズでのJWTトークン検証
- HTTPエンドポイントとの連携:価格変更がWebSocketプッシュ通知をトリガー
- Aliceシナリオ:リアルタイム価格フィード — 百万件の中の製品価格が変動したとき, 購読者が即座にプッシュ通知を受け取る
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) ライフサイクルステート
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プッシュアーキテクチャ
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のクロスオリジンハンドシェイクを自動的に処理します。
📖 まとめ
- WebSocketは永続的な双方向接続を提供し, ポーリングに代わる真のリアルタイムプッシュ通知を実現します
- ConnectionManagerは接続プールを管理し, プロダクトIDによるグループブロードキャストをサポートします
- JWT認証はハンドシェイクフェーズでクエリパラメータを使用して検証され, 失敗時は接続が即座に閉じられます (code 4001)
- HTTPエンドポイントで価格作成後,
manager.broadcast_to_product()を呼び出してプッシュ通知をトリガーします - 本番環境ではRedis Pub/Subでクロスプロセスブロードキャストが必要です。単一プロセスのメモリは複数ワーカーに適しません
📝 練習問題
- 基本問題 (難易度 ⭐):WebSocketエコーエンドポイントを作成し, クライアントがテキストを送信するとサーバーがそのまま返します。ブラウザの開発者ツールまたはwscatでテストします。ヒント:
@app.websocket("/ws/echo")+accept()+receive_text() - 応用問題 (難易度 ⭐⭐):複数クライアントが同じプロダクトの価格変更を購読できるConnectionManagerを実装します。新しい価格が作成されたとき, そのプロダクトを購読しているすべてのWebSocketクライアントにブロードキャストします。ヒント:
Dict[int, List[WebSocket]]+broadcast_to_product() - チャレンジ (難易度 ⭐⭐⭐):WebSocketにJWT認証を追加 (クエリパラメータでトークンを渡す)し, HTTP価格作成エンドポイントでWebSocketブロードキャストをトリガーして, 完全な「HTTP書き込み → WebSocketプッシュ」ワークフローを実装します。ヒント:
token: str = Query(...)+verify_ws_token()+ HTTPエンドポイントでmanager.broadcast_to_product()を呼び出し
---|



