IA-Ingenieria-MLOps · Skill 38252d5c

Arquitecto de Pipelines de Recomendacion

Diseño de pipeline de 6 etapas para el Feed "Para Ti" de CULTIVA IA — 2.022 skills, top-10 personalizadas por consultor

Cliente CULTIVA IA (interno)
Stack Python · FastAPI · async
Modo Online (request-time)
Items 2.022 skills
K 10
Estado Diseñado

Framework de 6 etapas (patrón For You Algorithm · xAI Apache 2.0)

1
Source
Candidatos del catálogo
2
Hydrator
Metadatos + historial
3
Filter
Eliminar inelegibles
4
Scorer
Multi-acción ponderado
5
Selector
Top-K + diversidad
6
SideEffect
Cache + analytics (async)

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).

Redis cache async/parallel ~2.200 cands
📌

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.

async gather Postgres Redis counters
🚫

3 · Filter (secuencial, barato→caro)

6 filtros — de ~2.200 a ~400

1
duplicate_filter
O(n)
2
inactive_filter
campo
3
served_recently_filter
Redis set
4
already_used_5x_filter
historial
5
nivel_eligibility_filter
lookup
6
acceso_filter
permisos
🎯

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.

candidatos aislados determinístico cacheable
🎉

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).

sort O(n log n) max 3/servicio fallback

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).

asyncio.create_task no await Redis Streams

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.

P(use)
×1.0
Usará la skill esta semana. Señal principal positiva.
P(share)
×0.5
La compartirá con un cliente. Amplificador de valor.
P(skip)
×−0.3
La ignorará. Penalización leve por irrelevancia.
P(archive)
×−0.5
La marcará como no relevante. Penalización fuerte.
Formula del score final
final_score = (1.0 × P_use) + (0.5 × P_share) + (-0.3 × P_skip) + (-0.5 × P_archive) + diversity_delta + biz_rules_delta

Scaffold Python generado (runnable · FastAPI async)

cultiva_feed/pipeline.py Python 3.12
# 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"

Feed Para Ti
user_id: marta.garcia · project: SEO · 2026-06-16 09:41 · 2.022 skills → 10 seleccionadas
#
Skill
P(use)
P(share)
Score
Razon
1
ai-seo
SEO · avanzado
0.91
0.72
1.18
contexto SEO
2
seo-audit
SEO · intermedio
0.88
0.65
1.10
contexto SEO
3
schema-markup
SEO · basico
0.79
0.81
1.04
tendencia
4
competitor-analysis
Marketing · avanzado
0.75
0.58
0.97
usuarios similares
5
content-strategy
Contenido · intermedio
0.71
0.63
0.90
usuarios similares
6
programmatic-seo
SEO · avanzado
0.68
0.71
0.89
popular global
7
copywriting
Contenido · basico
0.64
0.55
0.82
tendencia
8
cold-email
Ventas · intermedio
0.61
0.49
0.76
usuarios similares
9
analytics-tracking
Datos · basico
0.59
0.44
0.73
contexto SEO
10
paid-ads
Publicidad · intermedio
0.55
0.61
0.69
popular global

Presupuesto de latencia · modo online · objetivo < 150ms p99

Budget de Latencia por Etapa

Source 12ms
Hydrate 20ms
Filter 9ms
Scorer 65ms
Select 5ms
SFX ~0
margen 29ms
Source 12ms Hydrator 20ms (parallel) Filter 9ms Scorer 65ms (mayor costo) Selector 5ms SideEffects async ~0ms Margen 29ms Total: ~111ms p50 / 140ms p99

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

Scoring: unico vs multi-accion
Multi-accion (elegido)
Pesos ajustables sin reentrenar. Feed evolucionable sin datos historicos suficientes.
Score unico
Mas simple, pero requiere reentrenar para cambiar comportamiento.
Candidatos: aislados vs joint scoring
Candidatos aislados (elegido)
Deterministico, cacheable, compone bien con reranking de diversidad.
Joint scoring
Mas expresivo pero no-deterministico y no cacheable. Sin razon especifica, descartado.
Modo: online vs offline-batch
Online request-time (elegido)
El contexto del proyecto cambia por sesion. Necesitamos frescura inmediata.
Offline batch pre-computado
Menor latencia pero sin sensibilidad al proyecto activo. Mejor para notificaciones diarias.