Módulo 5: Materialized Views
Entrega del módulo 5: dashboard de blog con materialized views
¿Qué vas a construir y por qué?
Vas a implementar el sistema completo del dashboard interno de un blog usando 4 materialized views. El dashboard sirve métricas analíticas (top posts semanales, top comentaristas mensuales, tendencia de posts por mes, categorías más activas) que sin precómputo tardan ~17 segundos en total cargar. Con MVs bien diseñadas, baja a ~50 milisegundos. El refresh es agendado vía cron, las queries del frontend son lookups instantáneos sobre tablas indexadas, y el banner "última actualización" comunica el staleness al usuario.
Este proyecto integra todo lo aprendido en el módulo: la decisión de cuándo usar MV (cápsulas 02 y 07), creación con WITH NO DATA para deploys rápidos (cápsula 03), refresh CONCURRENTLY con unique index obligatorio (cápsula 04), indexes apropiados según las queries (cápsula 05), y arquitectura de "una MV por panel" (cápsula 06).
Por qué este proyecto es portfolio-worthy: entrega un PR realista con migrations Alembic, código FastAPI/SQLAlchemy 2.0 async, tests, cron job, y documentación arquitectónica (ADR). Un reviewer técnico senior puede evaluar tu nivel mirando: (a) si elegiste MV correctamente para cada panel; (b) si los unique indexes están desde el día 1; (c) si los indexes de queries cubren los patrones reales; (d) si el banner staleness está en los endpoints; (e) si el cron está protegido (lo que anticipamos para el módulo 6).
Conexión con el proyecto integrador (módulo 8): lo que entregas aquí es funcionalmente idéntico al componente "Materialized view top_posts_weekly" del proyecto final. El módulo 6 le agrega pg_try_advisory_lock para refresh seguro contra concurrencia. El módulo 8 lo integra al refactor del Blog API completo con FTS, partitioning, JSONB y CTEs recursivas.
Objetivo del proyecto
Al completar este proyecto, vas a tener:
- 4 materialized views creadas con migrations Alembic (1 por panel del dashboard).
- Un job de cron que refresca las MVs con frecuencias apropiadas.
- 4 endpoints FastAPI que consumen las MVs y devuelven
last_updated+staleness_minutes. - 1 endpoint agregador
GET /dashboardque devuelve los 4 paneles juntos. - Tests para los endpoints (mínimo: schema correcto, staleness presente).
- Benchmark documentado: latencia antes/después de cada panel.
ARCHITECTURE.mdcon decisiones justificadas (por qué MV y no Redis para cada panel).
Cómo encaja en lo que aprendiste
| Concepto del módulo | Dónde se usa en el proyecto |
|---|---|
| Cápsula 02: views vs materialized views | Decisión de MV vs query directo vs Redis para cada panel — documentada en ARCHITECTURE.md. |
| Cápsula 03: creación y refresh | Cada migration crea una MV con CREATE MATERIALIZED VIEW ... WITH NO DATA + el unique index obligatorio + primer refresh FULL. |
Cápsula 04: CONCURRENTLY vs FULL | Cron usa CONCURRENTLY siempre. Migration usa FULL para primer refresh post-WITH NO DATA. |
| Cápsula 05: indexes en MVs | Cada MV tiene indexes apropiados para las queries de su endpoint. |
| Cápsula 06: casos analíticos | Patrón "una MV por panel" + endpoint agregador + banner last_updated. |
| Cápsula 07: decision matrix | El ARCHITECTURE.md aplica la matriz para justificar cada decisión. |
Pensá el proyecto como el sistema mínimo viable de dashboard analítico interno. Cada pieza tiene propósito; nada está por relleno.
Especificaciones técnicas
Stack
- Lenguaje: Python 3.11+
- Framework principal: FastAPI 0.110+
- ORM: SQLAlchemy 2.0+ async
- DB: PostgreSQL 16+ (asumiendo
posts,views,comments,users,categoriesya existentes) - Migrations: Alembic
- Tests: pytest + httpx (para tests async)
- Cron: cron del SO (
*/60 * * * *) opg_cron(opcional)
Setup inicial
Asumiendo que ya tienes el proyecto del Blog API (de la guía #8) con tablas posts, views, comments, users, categories. Si no, puedes usar el setup de la cápsula 03 para tener tablas mínimas.
# Estructura del proyecto (extracto relevante)
app/
├── api/
│ └── dashboard.py # endpoints (este proyecto)
├── jobs/
│ └── refresh_dashboard.py # cron job (este proyecto)
├── models/
│ └── mv.py # mappings de las 4 MVs (este proyecto)
├── schemas/
│ └── dashboard.py # Pydantic schemas (este proyecto)
├── services/
│ └── mv_refresh.py # función refresh con whitelist (este proyecto)
└── alembic/
└── versions/
├── 042_mv_top_posts_weekly.py
├── 043_mv_top_commenters_monthly.py
├── 044_mv_posts_per_month.py
└── 045_mv_active_categories_weekly.py
ARCHITECTURE.md # documento de decisiones (este proyecto)
Funcionalidades obligatorias
1. MV mv_top_posts_weekly
Spec:
- 10 filas (top posts por views últimos 7 días).
- Columnas:
post_id,title,slug,view_count,computed_at. - Unique index sobre
post_id. - Refresh frecuencia: cada 1 hora.
Migration esperada (042_mv_top_posts_weekly.py):
"""create mv_top_posts_weekly
Revision ID: 042_mv_top_posts_weekly
"""
from alembic import op
def upgrade() -> None:
op.execute("""
CREATE MATERIALIZED VIEW mv_top_posts_weekly AS
SELECT
p.id AS post_id,
p.title,
p.slug,
count(v.id) AS view_count,
NOW() AS computed_at
FROM posts p
JOIN views v ON v.post_id = p.id
WHERE v.created_at > NOW() - INTERVAL '7 days'
AND p.published_at IS NOT NULL
GROUP BY p.id, p.title, p.slug
ORDER BY view_count DESC
LIMIT 10
WITH NO DATA;
""")
op.execute("""
CREATE UNIQUE INDEX idx_mv_top_posts_weekly_pk
ON mv_top_posts_weekly (post_id);
""")
# Primer refresh (FULL — necesario después de WITH NO DATA)
op.execute("REFRESH MATERIALIZED VIEW mv_top_posts_weekly;")
def downgrade() -> None:
op.execute("DROP MATERIALIZED VIEW IF EXISTS mv_top_posts_weekly;")
Endpoint: GET /dashboard/top-posts-weekly
Respuesta esperada:
{
"items": [
{"post_id": 142, "title": "Post X", "slug": "post-x", "view_count": 8421},
...
],
"last_updated": "2026-05-02T14:30:00Z",
"staleness_minutes": 47
}
2. MV mv_top_commenters_monthly
Spec:
- 10 filas (top comentaristas por # comentarios últimos 30 días).
- Columnas:
user_id,username,comment_count,computed_at. - Unique index sobre
user_id. - Refresh frecuencia: cada 1 hora.
Endpoint: GET /dashboard/top-commenters-monthly
3. MV mv_posts_per_month
Spec:
- 12 filas (un row por mes en últimos 12 meses).
- Columnas:
month,post_count,computed_at. - Unique index sobre
month. - Refresh frecuencia: cada 6 horas (cambia poco hora a hora).
Endpoint: GET /dashboard/posts-per-month
4. MV mv_active_categories_weekly
Spec:
- N filas (todas las categorías ranqueadas, top 10 mostradas).
- Columnas:
category_id,category_name,view_count,comment_count,activity_score,computed_at. - Unique index sobre
category_id. - Index adicional sobre
activity_score DESCpara elORDER BYdel endpoint. - Refresh frecuencia: cada 1 hora.
Endpoint: GET /dashboard/active-categories-weekly
5. Endpoint agregador
Spec: GET /dashboard devuelve los 4 paneles en una sola request, con overall_last_updated (el más viejo de los 4) y overall_staleness_minutes (el peor de los 4).
{
"top_posts_weekly": { ... },
"top_commenters_monthly": { ... },
"posts_per_month": { ... },
"active_categories_weekly": { ... },
"overall_last_updated": "2026-05-02T14:00:00Z",
"overall_staleness_minutes": 77
}
6. Cron de refresh
Spec: script Python que refresca las 4 MVs con CONCURRENTLY, registra duración y status, y se ejecuta vía cron del SO o pg_cron.
# jobs/refresh_dashboard.py
import asyncio
import logging
from app.db import async_session_factory
from services.mv_refresh import refresh_mv
logger = logging.getLogger(__name__)
MVS_HOURLY = [
"mv_top_posts_weekly",
"mv_top_commenters_monthly",
"mv_active_categories_weekly",
]
async def refresh_hourly_mvs() -> None:
"""Refresca las MVs que se actualizan cada hora."""
async with async_session_factory() as session:
for mv in MVS_HOURLY:
result = await refresh_mv(session, mv, concurrent=True)
if result["success"]:
logger.info(f"Refreshed {mv} in {result['duration_ms']}ms")
else:
logger.error(f"Failed {mv}: {result.get('error')}")
async def refresh_six_hourly_mvs() -> None:
"""Refresca las MVs que se actualizan cada 6 horas."""
async with async_session_factory() as session:
result = await refresh_mv(session, "mv_posts_per_month", concurrent=True)
if result["success"]:
logger.info(f"Refreshed mv_posts_per_month in {result['duration_ms']}ms")
else:
logger.error(f"Failed: {result.get('error')}")
if __name__ == "__main__":
import sys
if len(sys.argv) > 1 and sys.argv[1] == "six-hourly":
asyncio.run(refresh_six_hourly_mvs())
else:
asyncio.run(refresh_hourly_mvs())
Cron del SO:
# /etc/cron.d/dashboard-refresh
0 * * * * cd /app && python -m jobs.refresh_dashboard >> /var/log/refresh.log 2>&1
0 */6 * * * cd /app && python -m jobs.refresh_dashboard six-hourly >> /var/log/refresh.log 2>&1
Validaciones y manejo de errores
Qué debe validarse
- Cada MV tiene unique index — si falta,
REFRESH CONCURRENTLYfalla. - Endpoint maneja MV vacía (devuelve lista vacía con
staleness_minutes = 0). - Endpoint usa timezone-aware datetimes (no naive).
- Whitelist de MVs en
refresh_mv(evita SQL injection).
Errores que deben manejarse
- MV no populated (primer refresh fallido): el endpoint debe responder algo razonable (lista vacía + log de error), no crashear.
- Refresh falla por contención: el cron loggea el error sin abortar las siguientes MVs.
- Unique index missing: detectar en el código de
refresh_mvy devolver mensaje accionable ("MV needs UNIQUE INDEX for concurrent refresh").
Ejemplo de implementación mínima
Esqueleto funcional ejecutable. Te toca expandirlo para cumplir las funcionalidades obligatorias.
# app/api/dashboard.py — esqueleto mínimo
from datetime import datetime, timezone
from fastapi import APIRouter, Depends
from pydantic import BaseModel
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.deps import get_session
from app.models.mv import TopPostsWeekly # mapeado a mv_top_posts_weekly
router = APIRouter(prefix="/dashboard", tags=["dashboard"])
class TopPostItem(BaseModel):
post_id: int
title: str
slug: str
view_count: int
class TopPostsResponse(BaseModel):
items: list[TopPostItem]
last_updated: datetime
staleness_minutes: int
@router.get("/top-posts-weekly", response_model=TopPostsResponse)
async def get_top_posts_weekly(
session: AsyncSession = Depends(get_session),
) -> TopPostsResponse:
stmt = select(TopPostsWeekly).order_by(TopPostsWeekly.view_count.desc())
result = await session.execute(stmt)
rows = result.scalars().all()
if not rows:
now = datetime.now(timezone.utc)
return TopPostsResponse(items=[], last_updated=now, staleness_minutes=0)
last_updated = rows[0].computed_at
now = datetime.now(timezone.utc)
staleness = int((now - last_updated).total_seconds() / 60)
return TopPostsResponse(
items=[
TopPostItem(
post_id=r.post_id, title=r.title, slug=r.slug, view_count=r.view_count
)
for r in rows
],
last_updated=last_updated,
staleness_minutes=max(staleness, 0),
)
# TODO: implementar endpoints para los otros 3 paneles
# TODO: implementar endpoint agregador GET /dashboard
# app/services/mv_refresh.py — esqueleto
import time
from sqlalchemy import text
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.ext.asyncio import AsyncSession
ALLOWED_MVS = {
"mv_top_posts_weekly",
"mv_top_commenters_monthly",
"mv_posts_per_month",
"mv_active_categories_weekly",
}
async def refresh_mv(
session: AsyncSession,
mv_name: str,
*,
concurrent: bool = True,
) -> dict:
if mv_name not in ALLOWED_MVS:
return {"success": False, "error": f"MV '{mv_name}' not allowed"}
mode = "CONCURRENTLY " if concurrent else ""
sql = f"REFRESH MATERIALIZED VIEW {mode}{mv_name}"
started = time.monotonic()
try:
await session.execute(text(sql))
await session.commit()
duration_ms = (time.monotonic() - started) * 1000
return {"success": True, "duration_ms": round(duration_ms, 2)}
except SQLAlchemyError as e:
await session.rollback()
return {"success": False, "error": str(e)}
# TODO: detectar el error específico de "needs unique index"
# TODO: implementar fallback a FULL si MV no está populated (cápsula 04)
Este esqueleto:
- Es ejecutable (puedes correr el endpoint hoy).
- Muestra la estructura esperada (timezone-aware, schema con metadata, whitelist).
- NO incluye: los otros 3 endpoints, el cron, el endpoint agregador, los tests, la documentación. Eso es tu trabajo.
Rúbrica de evaluación (auto-verificación)
Total: 100 puntos. Aprobado: ≥70 puntos.
Funcionalidad (50 puntos)
- (10 pts) MV
mv_top_posts_weeklycreada con migration + unique index + primer refresh. - (10 pts) MV
mv_top_commenters_monthlycreada con migration + unique index + primer refresh. - (5 pts) MV
mv_posts_per_monthcreada con migration + unique index + primer refresh. - (10 pts) MV
mv_active_categories_weeklycreada con migration + unique index + index adicional sobreactivity_score DESC. - (5 pts) 4 endpoints individuales devuelven schema correcto +
last_updated+staleness_minutes. - (5 pts) Endpoint agregador
GET /dashboarddevuelve los 4 paneles +overall_*. - (5 pts) Cron job ejecuta refresh de las 4 MVs sin errores en condiciones normales.
Calidad de código (25 puntos)
- (5 pts) Type hints consistentes en código Python (Pydantic + Mapped en SQLAlchemy).
- (5 pts) Whitelist de MVs en
refresh_mv(no acepta nombre arbitrario). - (5 pts) Manejo de error razonable (try/except + rollback + log estructurado).
- (5 pts) Datetimes timezone-aware en todo el código.
- (5 pts) Tests de los endpoints (mínimo: status 200 + presencia de campos requeridos).
Documentación (15 puntos)
- (5 pts)
ARCHITECTURE.mdjustifica cada MV con criterios cuantitativos (latencia objetivo, staleness aceptable, frecuencia consulta). - (5 pts) Benchmark documentado: latencia antes/después de cada panel con
EXPLAIN ANALYZEantes yEXPLAIN ANALYZEdespués. - (5 pts) README explica cómo correr el cron (cron del SO o
pg_cron).
Extra credit (opcional, hasta +10 pts)
- (+3 pts) Endpoint que expone métricas del refresh (duración, last success, last failure) en
/admin/mv-status. - (+3 pts) Test de integración que verifica: refresh → cambio en datos base → refresh nuevamente → endpoint devuelve datos actualizados.
- (+2 pts) Implementación de fallback a
FULLcuando la MV no está populated (cápsula 04). - (+2 pts) Banner mock en frontend (HTML/JS sencillo) que consume el endpoint y muestra "Última actualización: hace X minutos".
Errores comunes en este proyecto
Error 1: olvidar el unique index en alguna migration
Síntoma: la migration corre OK, el primer refresh FULL también, pero el cron con CONCURRENTLY falla con cannot refresh materialized view ... concurrently.
Por qué pasa: te olvidaste CREATE UNIQUE INDEX en la migration. Es el error #1 en este proyecto.
Cómo corregir:
# en la migration, INMEDIATAMENTE después del CREATE MATERIALIZED VIEW:
op.execute("""
CREATE UNIQUE INDEX idx_<mv_name>_pk
ON <mv_name> (<unique_column>);
""")
Auditá tus 4 migrations antes de mergear.
Error 2: el endpoint falla con MV vacía (primer despliegue)
Síntoma: después de migration con WITH NO DATA, el endpoint devuelve 500 porque rows[0].computed_at falla con "list index out of range".
Por qué pasa: la MV está vacía (todavía no se refrescó), pero el endpoint asume que siempre hay al menos 1 fila.
Cómo corregir: chequea if not rows antes de acceder a rows[0]. El esqueleto que te di ya lo hace correctamente, pero asegúrate de mantenerlo al expandir.
Error 3: datetime naive vs timezone-aware
Síntoma: error can't compare offset-naive and offset-aware datetimes cuando calculas staleness_minutes.
Por qué pasa: datetime.utcnow() devuelve naive, computed_at desde la DB devuelve aware (porque la columna es TIMESTAMPTZ). Restarlos falla.
Cómo corregir: siempre usar datetime.now(timezone.utc):
from datetime import datetime, timezone
now = datetime.now(timezone.utc)
staleness = (now - last_updated).total_seconds() / 60
Nunca uses datetime.utcnow() (deprecado en Python 3.12 y propenso a este bug).
Error 4: cron sin whitelist de MVs
Síntoma: el cron acepta nombre arbitrario de MV vía argumento. Un error tipográfico ejecuta REFRESH MATERIALIZED VIEW some_typo que falla con SQL syntax error.
Por qué pasa: sin whitelist, el código es flexible pero permisivo. Más grave: si el nombre viene de input no controlado, es vulnerabilidad de SQL injection (los nombres de tabla no se pueden parametrizar).
Cómo corregir: whitelist explícita en ALLOWED_MVS. Solo nombres en esa lista son aceptados.
Error 5: refresh secuencial cuando podría ser paralelo
Síntoma: el cron tarda 17 segundos (suma de los 4 refreshes). En paralelo tardaría ~9s (el más lento, mv_active_categories_weekly).
Por qué pasa: implementación naïve loopea las MVs secuencialmente.
Cómo corregir: usar asyncio.gather con sessions independientes:
async def refresh_in_parallel() -> list[dict]:
async def refresh_one(mv_name):
async with async_session_factory() as session:
return await refresh_mv(session, mv_name, concurrent=True)
results = await asyncio.gather(*[refresh_one(mv) for mv in MVS_HOURLY])
return results
Trade-off: más conexiones a Postgres simultáneas. Si tu pool tiene cap bajo, considerar.
Error 6: olvidar commit() después del refresh
Síntoma: el refresh "parece" funcionar pero los datos no cambian desde la sesión que ejecutó el refresh.
Por qué pasa: SQLAlchemy 2.0 async opera en transacciones. Sin commit(), otros sessions no ven el cambio.
Cómo corregir: siempre await session.commit() después del refresh. Está en el esqueleto que te di — manténlo.
Error 7: el overall_last_updated se calcula al revés
Síntoma: el dashboard muestra overall_staleness_minutes: 0 aunque uno de los paneles tiene 60 minutos de staleness.
Por qué pasa: alguien usó max(last_updated) (el más reciente) en lugar de min(last_updated) (el más viejo).
Cómo corregir: lógicamente — el "overall last updated" es el momento desde el cual algún dato puede estar desactualizado, que es el más viejo de los timestamps:
overall_last_updated = min(all_last_updated)
overall_staleness_minutes = max(all_staleness)
¿Qué hacer si te atoras?
-
Setup de tablas no funciona: revisa la cápsula 03, sección "Setup de tablas y datos de prueba". Asegurate que
posts,views,comments,users,categoriesexisten. -
CREATE MATERIALIZED VIEWfalla: revisa la sintaxis exacta en la cápsula 03. El error usual es columna inexistente o un alias mal escrito. -
REFRESH CONCURRENTLYfalla: cápsula 04. Casi siempre es el unique index. -
El endpoint funciona pero devuelve datos viejos después del refresh: olvido del
commit()o estás consultando una sesión vieja. Revisá la cápsula 03. -
El cron no corre: si usas cron del SO, valida con
crontab -l. Para debug, ejecuta el script manualmente:python -m jobs.refresh_dashboard. -
Tests fallan en CI pero pasan local: asegúrate que la base de tests tiene las MVs creadas (correr migrations antes de tests).
Recursos para el proyecto
- PostgreSQL 16 — Materialized Views — referencia oficial completa.
- SQLAlchemy 2.0 — async patterns — para sessions y queries async.
- Alembic — autogenerate vs manual migrations — para entender por qué las MVs requieren
op.execute()(Alembic no autodetecta MVs). - Crunchy Data — Materialized Views Best Practices — checklist operacional.
- pganalyze — Monitoring Materialized View Refreshes — para el extra credit de métricas.
- pg_cron repo — opción alternativa al cron del SO.
- FastAPI — testing async endpoints — para los tests con httpx.
Lo que sigue
Lo que construiste acá es la base directa del componente "Materialized view top_posts_weekly" del proyecto final del módulo 8. Esa pieza la trasladás casi sin cambios al refactor del Blog API completo.
Pero antes hay un módulo intermedio crítico. Tu cron actual no protege contra refreshes concurrentes — si dos instancias arrancan al mismo tiempo (deploy doble, k8s reconciliando), el segundo se cuelga esperando al primero. El módulo 6 te enseña pg_try_advisory_lock — un lock no transaccional que permite que el segundo cron haga skip en lugar de esperar. Es la última pieza para que tu sistema sea production-grade. La transición narrativa: "tienes un cron que refresca cada hora. Pero ¿qué pasa si dos instancias arrancan al mismo tiempo? Para eso PostgreSQL tiene advisory locks — locks no transaccionales que reemplazan a Redis SETNX para coordinación distribuida."
Antes de avanzar al módulo 6, asegúrate de que tu proyecto:
- Pasa todos los tests.
- Cumple ≥70 puntos en la rúbrica.
- Tiene
ARCHITECTURE.mdjustificando cada MV con la decision matrix de la cápsula 07. - Tiene benchmark before/after documentado.
- Corre el cron exitosamente al menos 2 veces seguidas (para validar idempotencia).
Cuando esté listo, abre el módulo 6.
Módulo 5 — Advanced PostgreSQL for Backend Guide
Próximo módulo: Advisory Locks + Savepoints — coordinación distribuida sin Redis y rollbacks parciales con savepoints.