Módulo 2: ¿Cómo funcionan Embeddings?

Mini-Proyecto: API Client Robusto para Producción

Descripción del proyecto

Construirás un cliente de embeddings production-ready que implementa todos los patterns vistos en este módulo: batch processing, caching, retry logic con exponential backoff, rate limiting, normalización, y logging. Este cliente es reutilizable en proyectos reales.

Al completar este proyecto, tendrás una biblioteca Python sólida para generar embeddings a escala con todas las optimizaciones de producción.


Objetivos del proyecto

Funcionalidades:

  1. ✅ Generar embeddings (single y batch)
  2. ✅ Caching en disco (evitar re-generar)
  3. ✅ Retry logic (exponential backoff)
  4. ✅ Rate limiting (respetar RPM limits)
  5. ✅ Batch processing automático
  6. ✅ Normalización opcional
  7. ✅ Logging detallado
  8. ✅ CLI para testing

Estructura del proyecto

embeddings-client/
├── src/
│   ├── __init__.py
│   ├── client.py           # Cliente principal
│   ├── cache.py            # Sistema de cache
│   ├── rate_limiter.py     # Rate limiter
│   └── utils.py            # Utilidades (normalize, etc.)
├── tests/
│   └── test_client.py
├── cache/                  # Cache de embeddings (auto-generado)
├── .env
├── requirements.txt
└── README.md

Setup inicial

1. requirements.txt

openai==1.54.0
python-dotenv==1.0.0
numpy==1.26.4

2. .env

OPENAI_API_KEY=tu-api-key-aqui

3. Instalar dependencias

pip install -r requirements.txt

Implementación

Paso 1: utils.py (Utilidades)

"""
Utilidades para embeddings
"""
import numpy as np
from typing import List

def normalize_embedding(embedding: List[float]) -> List[float]:
    """
    Normalizar embedding a magnitud 1.0
    
    Args:
        embedding: Vector a normalizar
    
    Returns:
        Embedding normalizado
    """
    embedding = np.array(embedding)
    norm = np.linalg.norm(embedding)
    
    if norm == 0:
        return embedding.tolist()
    
    return (embedding / norm).tolist()

def cosine_similarity(emb_a: List[float], emb_b: List[float]) -> float:
    """
    Calcular cosine similarity entre dos embeddings
    
    Args:
        emb_a: Primer embedding
        emb_b: Segundo embedding
    
    Returns:
        Similarity score [0, 1]
    """
    a = np.array(emb_a)
    b = np.array(emb_b)
    
    return np.dot(a, b) / (np.linalg.norm(a) * np.linalg.norm(b))

Paso 2: cache.py (Sistema de cache)

"""
Sistema de cache para embeddings
"""
import json
import hashlib
from pathlib import Path
from typing import Optional, List

class EmbeddingCache:
    """Cache de embeddings en disco"""
    
    def __init__(self, cache_dir: str = "./cache"):
        """
        Inicializar cache
        
        Args:
            cache_dir: Directorio para almacenar cache
        """
        self.cache_dir = Path(cache_dir)
        self.cache_dir.mkdir(exist_ok=True)
    
    def _get_cache_key(self, text: str, model: str, dimensions: Optional[int]) -> str:
        """
        Generar key único para texto + config
        
        Args:
            text: Texto
            model: Modelo usado
            dimensions: Dimensiones (None = default)
        
        Returns:
            Hash MD5 como key
        """
        # Combinar texto + modelo + dimensions
        cache_string = f"{text}:{model}:{dimensions}"
        return hashlib.md5(cache_string.encode()).hexdigest()
    
    def get(self, text: str, model: str, dimensions: Optional[int] = None) -> Optional[List[float]]:
        """
        Obtener embedding del cache
        
        Returns:
            Embedding o None si no existe
        """
        cache_file = self.cache_dir / f"{self._get_cache_key(text, model, dimensions)}.json"
        
        if cache_file.exists():
            with open(cache_file, 'r') as f:
                data = json.load(f)
                return data['embedding']
        
        return None
    
    def set(self, text: str, model: str, dimensions: Optional[int], embedding: List[float]) -> None:
        """
        Guardar embedding en cache
        """
        cache_file = self.cache_dir / f"{self._get_cache_key(text, model, dimensions)}.json"
        
        with open(cache_file, 'w') as f:
            json.dump({
                'text': text[:100],  # Guardar snippet
                'model': model,
                'dimensions': dimensions,
                'embedding': embedding
            }, f)
    
    def clear(self) -> int:
        """
        Limpiar cache completo
        
        Returns:
            Cantidad de archivos eliminados
        """
        count = 0
        for cache_file in self.cache_dir.glob("*.json"):
            cache_file.unlink()
            count += 1
        return count

Paso 3: rate_limiter.py (Rate limiting)

"""
Rate limiter para respetar API limits
"""
import time
from collections import deque
from typing import Optional

class RateLimiter:
    """Rate limiter con sliding window"""
    
    def __init__(self, max_requests: int, window_seconds: int):
        """
        Inicializar rate limiter
        
        Args:
            max_requests: Máximo de requests en ventana
            window_seconds: Tamaño de ventana en segundos
        """
        self.max_requests = max_requests
        self.window_seconds = window_seconds
        self.requests = deque()
    
    def acquire(self) -> None:
        """
        Esperar si necesario para no exceder rate limit
        """
        now = time.time()
        
        # Remover requests fuera de ventana
        while self.requests and self.requests[0] < now - self.window_seconds:
            self.requests.popleft()
        
        # Si estamos en límite, esperar
        if len(self.requests) >= self.max_requests:
            sleep_time = self.requests[0] + self.window_seconds - now + 0.1
            if sleep_time > 0:
                time.sleep(sleep_time)
            self.requests.popleft()
        
        # Registrar request
        self.requests.append(time.time())
    
    def reset(self) -> None:
        """Resetear rate limiter"""
        self.requests.clear()

Paso 4: client.py (Cliente principal)

"""
Cliente robusto de embeddings para producción
"""
import logging
import time
from typing import List, Optional, Union
from openai import OpenAI, APIError, RateLimitError, APIConnectionError
from dotenv import load_dotenv
import os

from .cache import EmbeddingCache
from .rate_limiter import RateLimiter
from .utils import normalize_embedding

# Setup logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

class EmbeddingsClient:
    """Cliente production-ready de embeddings"""
    
    def __init__(
        self,
        model: str = "text-embedding-3-small",
        dimensions: Optional[int] = None,
        use_cache: bool = True,
        cache_dir: str = "./cache",
        max_retries: int = 3,
        max_rpm: int = 500,
        normalize: bool = False
    ):
        """
        Inicializar cliente
        
        Args:
            model: Modelo de embeddings
            dimensions: Dimensiones del embedding (None = default)
            use_cache: Usar cache en disco
            cache_dir: Directorio de cache
            max_retries: Máximo de reintentos
            max_rpm: Máximo requests per minute
            normalize: Normalizar embeddings a magnitud 1.0
        """
        # Setup OpenAI client
        load_dotenv()
        self.client = OpenAI(api_key=os.getenv("OPENAI_API_KEY"))
        
        # Config
        self.model = model
        self.dimensions = dimensions
        self.use_cache = use_cache
        self.max_retries = max_retries
        self.normalize = normalize
        
        # Cache
        if use_cache:
            self.cache = EmbeddingCache(cache_dir)
        
        # Rate limiter
        self.rate_limiter = RateLimiter(max_requests=max_rpm, window_seconds=60)
        
        # Stats
        self.stats = {
            'total_requests': 0,
            'cache_hits': 0,
            'api_calls': 0,
            'errors': 0
        }
        
        logger.info(f"EmbeddingsClient initialized (model={model}, cache={use_cache})")
    
    def embed(self, text: str) -> Optional[List[float]]:
        """
        Generar embedding para un texto
        
        Args:
            text: Texto a convertir
        
        Returns:
            Embedding o None si falla
        """
        self.stats['total_requests'] += 1
        
        # Intentar cache
        if self.use_cache:
            cached = self.cache.get(text, self.model, self.dimensions)
            if cached:
                self.stats['cache_hits'] += 1
                logger.debug(f"Cache hit: {text[:50]}...")
                return cached
        
        # Generar embedding
        embedding = self._generate_embedding_with_retry(text)
        
        if embedding is None:
            return None
        
        # Normalizar si es requerido
        if self.normalize:
            embedding = normalize_embedding(embedding)
        
        # Guardar en cache
        if self.use_cache and embedding:
            self.cache.set(text, self.model, self.dimensions, embedding)
        
        return embedding
    
    def embed_batch(self, texts: List[str], batch_size: int = 100) -> List[Optional[List[float]]]:
        """
        Generar embeddings para múltiples textos
        
        Args:
            texts: Lista de textos
            batch_size: Tamaño de batch
        
        Returns:
            Lista de embeddings
        """
        embeddings = []
        
        for i in range(0, len(texts), batch_size):
            batch = texts[i:i+batch_size]
            
            # Procesar batch
            batch_embeddings = self._process_batch(batch)
            embeddings.extend(batch_embeddings)
            
            logger.info(f"Batch {i//batch_size + 1}: {len(batch)} textos procesados")
        
        return embeddings
    
    def _process_batch(self, texts: List[str]) -> List[Optional[List[float]]]:
        """
        Procesar batch de textos (con cache)
        
        Args:
            texts: Lista de textos
        
        Returns:
            Lista de embeddings
        """
        embeddings = []
        texts_to_generate = []
        indices_to_generate = []
        
        # Check cache
        for idx, text in enumerate(texts):
            if self.use_cache:
                cached = self.cache.get(text, self.model, self.dimensions)
                if cached:
                    embeddings.append(cached)
                    self.stats['cache_hits'] += 1
                    continue
            
            # Needs generation
            embeddings.append(None)  # Placeholder
            texts_to_generate.append(text)
            indices_to_generate.append(idx)
        
        # Generate missing embeddings
        if texts_to_generate:
            generated = self._generate_embeddings_batch_with_retry(texts_to_generate)
            
            # Insert generated embeddings
            for idx, emb in zip(indices_to_generate, generated):
                if emb:
                    # Normalizar si requerido
                    if self.normalize:
                        emb = normalize_embedding(emb)
                    
                    embeddings[idx] = emb
                    
                    # Cache
                    if self.use_cache:
                        self.cache.set(texts[idx], self.model, self.dimensions, emb)
        
        return embeddings
    
    def _generate_embedding_with_retry(self, text: str) -> Optional[List[float]]:
        """
        Generar embedding con retry logic
        
        Args:
            text: Texto
        
        Returns:
            Embedding o None si falla
        """
        for attempt in range(self.max_retries):
            try:
                # Rate limiting
                self.rate_limiter.acquire()
                
                # API call
                response = self.client.embeddings.create(
                    model=self.model,
                    input=text,
                    dimensions=self.dimensions
                )
                
                self.stats['api_calls'] += 1
                return response.data[0].embedding
            
            except RateLimitError:
                wait_time = 2 ** attempt
                logger.warning(f"Rate limit. Retry {attempt+1}/{self.max_retries} en {wait_time}s")
                time.sleep(wait_time)
            
            except (APIError, APIConnectionError) as e:
                logger.error(f"API error: {e}")
                if attempt < self.max_retries - 1:
                    time.sleep(1)
        
        # Failed
        self.stats['errors'] += 1
        logger.error(f"Failed to generate embedding para: {text[:50]}...")
        return None
    
    def _generate_embeddings_batch_with_retry(self, texts: List[str]) -> List[Optional[List[float]]]:
        """
        Generar batch de embeddings con retry
        
        Args:
            texts: Lista de textos
        
        Returns:
            Lista de embeddings
        """
        for attempt in range(self.max_retries):
            try:
                # Rate limiting
                self.rate_limiter.acquire()
                
                # API call (batch)
                response = self.client.embeddings.create(
                    model=self.model,
                    input=texts,
                    dimensions=self.dimensions
                )
                
                self.stats['api_calls'] += 1
                return [item.embedding for item in response.data]
            
            except RateLimitError:
                wait_time = 2 ** attempt
                logger.warning(f"Rate limit. Retry {attempt+1}/{self.max_retries} en {wait_time}s")
                time.sleep(wait_time)
            
            except (APIError, APIConnectionError) as e:
                logger.error(f"API error: {e}")
                if attempt < self.max_retries - 1:
                    time.sleep(1)
        
        # Failed
        self.stats['errors'] += len(texts)
        logger.error(f"Failed to generate batch of {len(texts)} embeddings")
        return [None] * len(texts)
    
    def get_stats(self) -> dict:
        """
        Obtener estadísticas del cliente
        
        Returns:
            Dict con stats
        """
        cache_hit_rate = 0
        if self.stats['total_requests'] > 0:
            cache_hit_rate = self.stats['cache_hits'] / self.stats['total_requests']
        
        return {
            **self.stats,
            'cache_hit_rate': f"{cache_hit_rate:.2%}"
        }
    
    def clear_cache(self) -> int:
        """
        Limpiar cache
        
        Returns:
            Cantidad de archivos eliminados
        """
        if self.use_cache:
            return self.cache.clear()
        return 0

Testing

Crear test_client.py

"""
Tests del cliente de embeddings
"""
from src.client import EmbeddingsClient
from src.utils import cosine_similarity

def test_single_embedding():
    """Test: Generar 1 embedding"""
    client = EmbeddingsClient(use_cache=False)
    
    embedding = client.embed("Python es un lenguaje")
    
    assert embedding is not None
    assert len(embedding) == 1536  # Default dimensions
    print("✅ Test single embedding passed")

def test_batch_embedding():
    """Test: Generar batch de embeddings"""
    client = EmbeddingsClient(use_cache=False)
    
    texts = ["Python", "JavaScript", "Go"]
    embeddings = client.embed_batch(texts)
    
    assert len(embeddings) == 3
    assert all(emb is not None for emb in embeddings)
    print("✅ Test batch embedding passed")

def test_cache():
    """Test: Cache funciona"""
    client = EmbeddingsClient(use_cache=True)
    
    text = "Test caching"
    
    # Primera llamada (API)
    emb1 = client.embed(text)
    api_calls_1 = client.stats['api_calls']
    
    # Segunda llamada (cache)
    emb2 = client.embed(text)
    api_calls_2 = client.stats['api_calls']
    
    assert emb1 == emb2
    assert api_calls_2 == api_calls_1  # No new API call
    assert client.stats['cache_hits'] == 1
    
    print("✅ Test cache passed")
    
    # Cleanup
    client.clear_cache()

def test_normalization():
    """Test: Normalización funciona"""
    import numpy as np
    
    client = EmbeddingsClient(normalize=True, use_cache=False)
    
    embedding = client.embed("Python es popular")
    magnitude = np.linalg.norm(embedding)
    
    assert 0.99 <= magnitude <= 1.01  # Magnitude ~1.0
    print(f"✅ Test normalization passed (magnitude={magnitude:.4f})")

def test_reduced_dimensions():
    """Test: Dimensiones reducidas"""
    client = EmbeddingsClient(dimensions=512, use_cache=False)
    
    embedding = client.embed("Python")
    
    assert len(embedding) == 512
    print("✅ Test reduced dimensions passed")

if __name__ == "__main__":
    print("Running tests...\n")
    
    test_single_embedding()
    test_batch_embedding()
    test_cache()
    test_normalization()
    test_reduced_dimensions()
    
    print("\n✅ Todos los tests pasaron!")

CLI de testing

Crear main.py

"""
CLI para testing del cliente
"""
from src.client import EmbeddingsClient
from src.utils import cosine_similarity

def main():
    """Demo interactivo"""
    print("=== Embeddings Client Demo ===\n")
    
    # Inicializar cliente
    client = EmbeddingsClient(
        model="text-embedding-3-small",
        dimensions=512,  # Reduced
        use_cache=True,
        normalize=True,
        max_rpm=100
    )
    
    # Test 1: Single embedding
    print("Test 1: Single embedding")
    embedding = client.embed("Python es un lenguaje de programación")
    print(f"✅ Embedding generado: {len(embedding)} dims\n")
    
    # Test 2: Similaridad
    print("Test 2: Similaridad semántica")
    emb_a = client.embed("Python es popular")
    emb_b = client.embed("Python es muy usado")
    emb_c = client.embed("Gato duerme en sofá")
    
    sim_ab = cosine_similarity(emb_a, emb_b)
    sim_ac = cosine_similarity(emb_a, emb_c)
    
    print(f"'Python es popular' vs 'Python es muy usado': {sim_ab:.4f}")
    print(f"'Python es popular' vs 'Gato duerme en sofá': {sim_ac:.4f}\n")
    
    # Test 3: Batch
    print("Test 3: Batch processing")
    texts = [f"Documento {i}" for i in range(10)]
    embeddings = client.embed_batch(texts, batch_size=5)
    print(f"✅ {len(embeddings)} embeddings generados\n")
    
    # Stats
    print("Stats:")
    stats = client.get_stats()
    for key, value in stats.items():
        print(f"  {key}: {value}")

if __name__ == "__main__":
    main()

Ejecutar:

python main.py

Output esperado:

=== Embeddings Client Demo ===

Test 1: Single embedding
✅ Embedding generado: 512 dims

Test 2: Similaridad semántica
'Python es popular' vs 'Python es muy usado': 0.8721
'Python es popular' vs 'Gato duerme en sofá': 0.4532

Test 3: Batch processing
✅ 10 embeddings generados

Stats:
  total_requests: 13
  cache_hits: 0
  api_calls: 3
  errors: 0
  cache_hit_rate: 0.00%

Validación del proyecto

Checklist:

  • Cliente genera embeddings single y batch ✅
  • Cache funciona (segunda llamada no llama API) ✅
  • Retry logic maneja rate limits ✅
  • Rate limiter respeta RPM ✅
  • Normalización funciona (magnitude ~1.0) ✅
  • Dimensiones reducidas funcionan (512 dims) ✅
  • Logging informa progreso ✅
  • Tests pasan sin errores ✅

Extensiones opcionales

1. Add async support (para alto volumen):

import asyncio
from openai import AsyncOpenAI

class AsyncEmbeddingsClient:
    """Version async del cliente"""
    
    async def embed(self, text: str):
        # Implementa versión async
        pass

2. Add database cache (Redis/PostgreSQL):

# En lugar de JSON files, usa Redis:
import redis

class RedisCache:
    def __init__(self):
        self.redis = redis.Redis(host='localhost', port=6379)
    
    def get(self, key):
        # Implementa
        pass

Resumen del proyecto

Qué implementaste:

  • Cliente robusto: Production-ready embeddings client
  • Caching: Disk-based cache (JSON)
  • Retry logic: Exponential backoff automático
  • Rate limiting: Sliding window rate limiter
  • Batch processing: Automático con cache integration
  • Normalización: Opcional (magnitud 1.0)
  • Logging: Detallado para monitoring
  • Tests: Suite completa de validación

Patrones aplicados:

  1. Exponential backoff (manejo de rate limits)
  2. Cache-aside pattern (check cache first)
  3. Batch processing (reduce API calls)
  4. Dependency injection (config flexible)

Conclusión del Módulo 2

Qué aprendiste en el módulo:

Arquitectura:

  • ✅ Transformer encoders (self-attention, multi-head attention)
  • ✅ Tokenización (BPE, WordPiece, tiktoken)
  • ✅ Contextualización (embeddings dinámicos)
  • ✅ Pooling strategies (mean, CLS, max)
  • ✅ Normalización (L2 norm, magnitud unitaria)

API de producción:

  • ✅ OpenAI API avanzada (batch, dimensions, retry)
  • ✅ Patterns de producción (caching, rate limiting)
  • ✅ Cliente robusto (proyecto completo)

Líneas de código generadas: ~1,200 (production-ready)


Siguiente módulo

Módulo 3: Modelos de Embeddings Comparison

Aprenderás:

  • OpenAI vs Sentence-BERT vs BGE
  • Open-source embeddings (local)
  • Benchmarking (MTEB)
  • Domain-specific embeddings
  • Multi-language embeddings
  • Cuándo usar qué modelo

De arquitectura a comparación de modelos.


Módulo 2 completadoPipeline completo: Tokens → Contextualización → Pooling → Normalization → Producción