Framework de 6 etapas (patrón For You Algorithm · xAI Apache 2.0)
Detalle por etapa
1 · Source (paralelo)
2 fuentes · 2.022 candidatos totales
skills_catalog_source — todas las skills activas del arsenal (slug, servicio, tipo, nivel, bucket, precio_eur). Paginado desde Redis con TTL 5 min.
skills_collab_source — skills usadas por consultores con perfil similar al usuario activo (user-user CF, top-200).
2 · Hydrator (paralelo)
3 hidratadores independientes
usage_hydrator — historial de uso del user: cuántas veces usó cada skill, última fecha, rating.
context_hydrator — etiquetas del proyecto activo vs. etiquetas de la skill (overlap score).
popularity_hydrator — conteo de usos globales de la skill en los últimos 30d.
3 · Filter (secuencial, barato→caro)
6 filtros — de ~2.200 a ~400
4 · Scorer chain
3 scorers secuenciales · multi-accion
multi_action_scorer — predice P(use), P(share), P(skip), P(archive). Combina con pesos configurables sin reentrenar.
diversity_scorer — penaliza skills del mismo servicio ya en el batch (–0.15 por servicio repetido).
business_rules_scorer — boost +0.2 a skills en promocion; –0.5 a skills bloqueadas por admin.
5 · Selector
Sort desc · top-10 · mix estratificado
Ordena por final_score descendente. Aplica stratified mix: max 3 skills del mismo servicio en el top-10 para garantizar variedad.
Si el batch tiene < 10 skills tras filtros, rellena con las más populares globalmente (fallback seguro).
6 · SideEffect (fire-and-forget)
Nunca bloquea la respuesta
impression_logger — emite evento skill.impressed a la cola de eventos (Redis Streams).
served_ids_cache — escribe los IDs servidos en Redis set (TTL 24h) para el filtro served_recently.
analytics_counter — incrementa contadores de impresión en Postgres (batch flush cada 60s).
Scoring multi-accion — pesos configurables sin reentrenar
Combiner de Acciones
Se predicen 4 probabilidades independientes. El score final es la suma ponderada. Para cambiar el comportamiento del feed, basta cambiar los pesos — sin reentrenar ningun modelo.
Scaffold Python generado (runnable · FastAPI async)
# cultiva_feed/pipeline.py # Feed "Para Ti" — CULTIVA IA · Arquitecto de Pipelines de Recomendacion # Patron: Source → Hydrator → Filter → Scorer → Selector → SideEffect from __future__ import annotations import asyncio from dataclasses import dataclass, field from typing import Protocol, Sequence from fastapi import FastAPI, Depends # ── Tipos base ─────────────────────────────────────────────────────────────── @dataclass class SkillCandidate: slug: str servicio: str tipo: str nivel: str bucket: str precio_eur: float acceso: str = "free" # hydrated fields usage_count: int = 0 context_overlap: float = 0.0 global_popularity: int = 0 # scores p_use: float = 0.0 p_share: float = 0.0 p_skip: float = 0.0 p_archive: float = 0.0 final_score: float = 0.0 @dataclass class FeedContext: user_id: str project_context: str # e.g. "SEO", "CRO", "Automatizacion" session_ts: float user_nivel: str = "intermedio" plan: str = "premium" # ── Protocolos (interfaces) ─────────────────────────────────────────────────── class Source(Protocol): async def fetch(self, ctx: FeedContext) -> list[SkillCandidate]: ... class Hydrator(Protocol): async def hydrate(self, cands: list[SkillCandidate], ctx: FeedContext) -> None: ... class Filter(Protocol): def apply(self, cands: list[SkillCandidate], ctx: FeedContext) -> list[SkillCandidate]: ... class Scorer(Protocol): def score(self, cands: list[SkillCandidate], ctx: FeedContext) -> None: ... class SideEffect(Protocol): async def run(self, top_k: list[SkillCandidate], ctx: FeedContext) -> None: ... # ── Pipeline principal ──────────────────────────────────────────────────────── @dataclass class FeedPipeline: sources: list[Source] hydrators: list[Hydrator] filters: list[Filter] scorers: list[Scorer] side_effects: list[SideEffect] top_k: int = 10 max_per_svc: int = 3 # diversidad: max skills por servicio async def run(self, ctx: FeedContext) -> list[SkillCandidate]: # 1. SOURCE — paralelo batches = await asyncio.gather(*[s.fetch(ctx) for s in self.sources]) cands = dedupe(c for b in batches for c in b) # 2. HYDRATOR — paralelo await asyncio.gather(*[h.hydrate(cands, ctx) for h in self.hydrators]) # 3. FILTER — secuencial (cheap → expensive) for f in self.filters: cands = f.apply(cands, ctx) # 4. SCORER — secuencial (cada scorer ve los anteriores) for sc in self.scorers: sc.score(cands, ctx) # 5. SELECTOR — sort + top-K + mix estratificado top_k = self._select(cands) # 6. SIDE EFFECTS — fire-and-forget, nunca bloquea for se in self.side_effects: asyncio.create_task(se.run(top_k, ctx)) # no await return top_k def _select(self, cands: list[SkillCandidate]) -> list[SkillCandidate]: sorted_c = sorted(cands, key=lambda c: c.final_score, reverse=True) result, svc_count = [], {} for c in sorted_c: if len(result) == self.top_k: break if svc_count.get(c.servicio, 0) < self.max_per_svc: result.append(c) svc_count[c.servicio] = svc_count.get(c.servicio, 0) + 1 return result # ── Scorer multi-accion ─────────────────────────────────────────────────────── @dataclass class MultiActionScorer: w_use: float = 1.0 w_share: float = 0.5 w_skip: float = -0.3 w_archive: float = -0.5 def score(self, cands: list[SkillCandidate], ctx: FeedContext) -> None: for c in cands: c.p_use = _predict_use(c, ctx) c.p_share = _predict_share(c, ctx) c.p_skip = _predict_skip(c, ctx) c.p_archive = _predict_archive(c, ctx) c.final_score = ( self.w_use * c.p_use + self.w_share * c.p_share + self.w_skip * c.p_skip + self.w_archive * c.p_archive ) # ── FastAPI endpoint ────────────────────────────────────────────────────────── app = FastAPI(title="CULTIVA Feed API") @app.get("/feed/{user_id}") async def get_feed(user_id: str, project: str = "general"): ctx = FeedContext(user_id=user_id, project_context=project, session_ts=time()) top_k = await pipeline.run(ctx) return [asdict(c) for c in top_k]
Ejemplo de output — Top-10 para consultor "marta.garcia" · proyecto "SEO"
Presupuesto de latencia · modo online · objetivo < 150ms p99
Budget de Latencia por Etapa
Optimizacion: Las 5 etapas antes del Scorer reducen los candidatos de ~2.200 a ~400, lo que hace que el paso mas caro (Scorer) procese 5.5x menos items. El filtrado barato primero es la optimizacion mas importante del pipeline.
Decisiones de diseno — trade-offs elegidos