WebSocket — Comunicação Bidirecional em Tempo Real
HTTP é como enviar uma carta—vai em uma direção e volta; WebSocket é como fazer uma ligação telefônica—ambas as partes podem falar a qualquer momento sem precisar desligar e ligar novamente.
1. O Que Você Vai Aprender
- Básico de WebSocket: Decoradores
@app.websockete Ciclo de Vida (connect/receive/disconnect) - Padrão Connection Manager: Design de Classe
ConnectionManagere Mecanismo de Broadcasting - WebSocket Autenticado: Verificar o token JWT durante a fase de handshake
- Integra com endpoints HTTP: Mudanças de preço acionam notificações push WebSocket
- Cenário da Alice: Feed de Preço em Tempo Real—Quando o preço de um produto muda entre milhões de itens, assinantes recebem imediatamente uma notificação push
2. A História Real da Alice
(1) Problema: Mudanças de preço só podem ser recuperadas via polling
O frontend do Bob faz polling à API do PriceTracker a cada 5 segundos para verificar mudanças de preço, mas 99% dos preços de milhões de produtos permanecem inalterados em qualquer período de 5 segundos, significando que 99% das requisições são desperdiçadas. Para piorar, podem levar até 5 segundos para que mudanças de preço sejam exibidas, levando a reclamações de clientes de que "os preços não são em tempo real."
(2) Solução WebSocket
WebSocket estabelece uma conexão persistente e bidirecional e envia ativamente mudanças de preço do servidor para o frontend quando ocorrem, eliminando a necessidade de polling—desperdício zero, latência zero.
@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) Resultado
O número de requisições de polling caiu de 200 por segundo para 0, o atraso na atualização de preços diminuiu de 5 segundos para 50 ms, o frontend do Bob não desperdiça mais cotas de API, e a satisfação do cliente melhorou significativamente.
3. Básico de WebSocket
(1) Estados do Ciclo de Vida
stateDiagram-v2
[*] --> CONNECTING: Cliente inicia
CONNECTING --> CONNECTED: accept()
CONNECTED --> RECEIVING: receive()
RECEIVING --> CONNECTED: send()
CONNECTED --> CLOSING: close() / disconnect
CLOSING --> CLOSED: Conexão fechada
CLOSED --> [*]
(1) ▶ Exemplo: Um Endpoint WebSocket Mínimo
from fastapi import FastAPI, WebSocket
app = FastAPI()
@app.websocket("/ws/echo")
async def websocket_echo(websocket: WebSocket):
await websocket.accept() # Aceitar conexão
try:
while True:
data = await websocket.receive_text()
await websocket.send_text(f"Echo: {data}")
except Exception:
await websocket.close()
Saída:
# Função definida com sucesso
(2) ▶ Exemplo: WebSocket com Parâmetros de Caminho
@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()
# Retornar eco com contexto do produto
await websocket.send_json({
"product_id": product_id,
"price": data.get("price"),
"currency": data.get("currency", "USD"),
})
except Exception:
await websocket.close()
Saída:
# Função definida com sucesso
4. Padrão Connection Manager
(1) Design do ConnectionManager
(1) ▶ Exemplo: Implementação do ConnectionManager
from fastapi import FastAPI, WebSocket
from typing import Dict, List
import json
class ConnectionManager:
def __init__(self):
# Mapeamento: product_id -> lista de conexões ativas
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)
# Limpar conexões mortas
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()
Saída:
# Função definida com sucesso
(2) Arquitetura de Push WebSocket
flowchart TD
Update[Atualização de Preço via HTTP] --> Handler[API Handler]
Handler --> Manager[ConnectionManager]
Manager --> WS1[Cliente WebSocket 1]
Manager --> WS2[Cliente WebSocket 2]
Manager --> WSN[Cliente WebSocket N]
subgraph Subscribers
WS1
WS2
WSN
end
(2) ▶ Exemplo: Usando ConnectionManager com um Endpoint WebSocket
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:
# Manter conexão viva, receber quaisquer mensagens do cliente
data = await websocket.receive_text()
except Exception:
manager.disconnect(websocket, product_id)
Saída:
# Função definida com sucesso
5. WebSocket Autenticado
(1) Verificando o JWT Durante a Fase de Handshake
WebSocket não possui um mecanismo padrão de cabeçalho; tokens são passados via parâmetros de consulta.
(1) ▶ Exemplo: Autenticação JWT WebSocket
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"),
):
# Verificar token antes de aceitar conexão
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
Saída:
# Função definida com sucesso
| Método de Autenticação | Implementação | Vantagens | Desvantagens |
|---|---|---|---|
| Parâmetro de Consulta | ?token=xxx |
Simples | Token aparece em logs de URL |
| Primeira mensagem | Enviar token após conectar | Não exposto na URL | Uma viagem adicional |
| Sec-WebSocket-Protocol | Token transmitido via subprotocolo | Não exposto na URL | Uso não padrão |
6. Colaboração HTTP e WebSocket
(1) Mudanças de Preço Acionam Notificações Push
(1) ▶ Exemplo: Endpoint HTTP Aciona um Broadcast WebSocket
from pydantic import BaseModel, Field
from datetime import datetime
class PriceUpdate(BaseModel):
product_id: int = Field(gt=0)
price: float = Field(gt=0, description="Novo preço em 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)
# Acionar broadcast WebSocket após criação bem-sucedida do preço
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
Saída:
# Função definida com sucesso
(2) ▶ Exemplo: Feed de Preço Global (Assinar Todas as Mudanças)
@app.websocket("/ws/prices/stream")
async def global_price_stream(websocket: WebSocket):
await websocket.accept()
# Adicionar a uma lista global de assinantes
manager.global_connections.append(websocket)
try:
while True:
# Receber heartbeat/ping do cliente
data = await websocket.receive_text()
if data == "ping":
await websocket.send_text("pong")
except Exception:
manager.global_connections.remove(websocket)
Saída:
# Função definida com sucesso
❓ Perguntas Frequentes
MessageModel.model_validate(data); se a validação falhar, envie uma mensagem de erro ao cliente.📖 Resumo
- WebSocket fornece uma conexão persistente e bidirecional, substituindo polling para habilitar notificações push verdadeiramente em tempo real
- ConnectionManager gerencia o pool de conexões e suporta broadcasting agrupado por ID de produto
- Autenticação JWT é verificada durante a fase de handshake usando parâmetros de consulta; se falhar, a conexão é fechada imediatamente (código 4001).
- Após criar um preço em um endpoint HTTP, chame
manager.broadcast_to_product()para acionar uma notificação push - O ambiente de produção requer Redis Pub/Sub para habilitar broadcasting entre processos; memória de processo único não é adequada para múltiplos workers.
📝 Exercícios
- Exercício Básico (Dificuldade: ⭐): Crie um endpoint echo WebSocket onde o cliente envia texto e o servidor retorna exatamente como recebido. Teste usando as ferramentas de desenvolvedor do seu navegador ou wscat. Dica:
@app.websocket("/ws/echo")+accept()+receive_text() - Problema Avançado (Dificuldade ⭐⭐): Implemente um ConnectionManager que suporta múltiplos clientes assinando mudanças de preço do mesmo produto. Quando um novo preço é criado, faça broadcast para todos os clientes WebSocket assinados naquele produto. Dica:
Dict[int, List[WebSocket]]+broadcast_to_product() - Desafio (Dificuldade: ⭐⭐⭐): Adicione autenticação JWT ao WebSocket (passando um token nos parâmetros de consulta), e acione um broadcast WebSocket no endpoint HTTP de criação de preço para implementar um fluxo completo de "escrita HTTP → push WebSocket". Dica:
token: str = Query(...)+verify_ws_token()+ endpoint HTTP chamandomanager.broadcast_to_product()
---|



