Módulo 7: Casos de Uso
8. Proyecto: Use Case Selector
Descripción
En este proyecto construyes un Use Case Selector: un router inteligente que recibe un input (PDF, imagen, audio, video, texto), detecta automáticamente su tipo, y aplica el pipeline correcto. Es el "cerebro" que conecta todos los patrones de diseño que aprendiste en este módulo en un solo sistema.
Por qué importa: En un sistema multimodal real, el usuario no dice "aplica el pipeline de Document Q&A". Sube un archivo y espera que el sistema haga lo correcto. El Use Case Selector es esa capa de inteligencia: recibe cualquier input, decide qué hacer con él, y ejecuta el pipeline apropiado. Es el componente más crítico del Document Analyzer que construirás en el Módulo 8.
Lo que vas a construir:
┌──────────────┐ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ Input │────▶│ Detector │────▶│ Pipeline │────▶│ Router │
│ (cualquier) │ │ de tipo │ │ Registry │ │ (ejecutar) │
└──────────────┘ └──────────────┘ └──────────────┘ └──────┬───────┘
│
▼
┌──────────────┐
│ Resultado │
│ + tracking │
└──────────────┘
Especificaciones
Input
- Ruta a archivo — PDF, imagen, audio, video, texto
- Pregunta opcional — Si el usuario quiere Q&A sobre el archivo
- Configuración — Modelo preferido, presupuesto máximo, calidad deseada
Output
- Tipo detectado — Qué tipo de input es
- Pipeline seleccionado — Qué pipeline se aplicó
- Resultado — Output del pipeline ejecutado
- Metadatos — Costo, duración, modelo usado
Componentes
| # | Componente | Responsabilidad |
|---|---|---|
| 1 | Config | Configuración global y por pipeline |
| 2 | Type Detector | Detectar tipo de input |
| 3 | Pipeline Registry | Registrar y gestionar pipelines disponibles |
| 4 | Router | Seleccionar y ejecutar pipeline |
| 5 | Cost Tracker | Rastrear costos por operación |
| 6 | Main | Orquestar todo |
Paso 1: Configuración
from dataclasses import dataclass, field
@dataclass
class PipelineConfig:
model: str = "gpt-4o-mini"
max_tokens: int = 500
temperature: float = 0.2
max_cost_per_request: float = 0.50
enable_cache: bool = True
enable_logging: bool = True
@dataclass
class UseCaseSelectorConfig:
default_pipeline_config: PipelineConfig = field(default_factory=PipelineConfig)
supported_types: set = field(default_factory=lambda: {
"document", "image", "audio", "video", "text"
})
max_file_size_mb: float = 100.0
fallback_enabled: bool = True
type_extensions: dict = field(default_factory=lambda: {
".pdf": "document",
".doc": "document",
".docx": "document",
".png": "image",
".jpg": "image",
".jpeg": "image",
".gif": "image",
".webp": "image",
".bmp": "image",
".mp3": "audio",
".wav": "audio",
".m4a": "audio",
".ogg": "audio",
".flac": "audio",
".mp4": "video",
".avi": "video",
".mov": "video",
".mkv": "video",
".webm": "video",
".txt": "text",
".md": "text",
".csv": "text",
".json": "text",
})
Paso 2: Detector de Tipo
from pathlib import Path
import mimetypes
class TypeDetector:
def __init__(self, config: UseCaseSelectorConfig):
self.config = config
def detect(self, path: str) -> dict:
p = Path(path)
if not p.exists():
return {
"type": "unknown",
"error": "file_not_found",
"path": path
}
ext = p.suffix.lower()
size_mb = p.stat().st_size / (1024 * 1024)
if size_mb > self.config.max_file_size_mb:
return {
"type": "unknown",
"error": "file_too_large",
"size_mb": round(size_mb, 2),
"max_mb": self.config.max_file_size_mb,
"path": path
}
detected_type = self.config.type_extensions.get(ext)
if not detected_type:
mime_type, _ = mimetypes.guess_type(path)
if mime_type:
detected_type = self._mime_to_type(mime_type)
if not detected_type:
detected_type = "unknown"
return {
"type": detected_type,
"extension": ext,
"size_mb": round(size_mb, 2),
"path": path,
"mime_type": mimetypes.guess_type(path)[0]
}
def _mime_to_type(self, mime: str) -> str | None:
if mime.startswith("image/"):
return "image"
if mime.startswith("audio/"):
return "audio"
if mime.startswith("video/"):
return "video"
if mime.startswith("text/"):
return "text"
if mime == "application/pdf":
return "document"
return None
def detect_multiple(self, paths: list[str]) -> list[dict]:
return [self.detect(path) for path in paths]
def validate(self, path: str) -> dict:
detection = self.detect(path)
if detection["type"] == "unknown":
return {
"valid": False,
"detection": detection,
"message": detection.get("error", "Tipo no soportado")
}
if detection["type"] not in self.config.supported_types:
return {
"valid": False,
"detection": detection,
"message": f"Tipo '{detection['type']}' no está habilitado"
}
return {
"valid": True,
"detection": detection
}
Paso 3: Pipeline Registry
from typing import Callable, Any
class Pipeline:
def __init__(
self,
name: str,
input_type: str,
handler: Callable,
description: str = "",
supports_question: bool = False,
estimated_cost_per_call: float = 0.01
):
self.name = name
self.input_type = input_type
self.handler = handler
self.description = description
self.supports_question = supports_question
self.estimated_cost_per_call = estimated_cost_per_call
class PipelineRegistry:
def __init__(self):
self.pipelines: dict[str, list[Pipeline]] = {}
def register(self, pipeline: Pipeline) -> None:
if pipeline.input_type not in self.pipelines:
self.pipelines[pipeline.input_type] = []
self.pipelines[pipeline.input_type].append(pipeline)
def get_pipeline(self, input_type: str, has_question: bool = False) -> Pipeline | None:
candidates = self.pipelines.get(input_type, [])
if not candidates:
return None
if has_question:
qa_pipelines = [p for p in candidates if p.supports_question]
if qa_pipelines:
return qa_pipelines[0]
return candidates[0]
def list_pipelines(self) -> dict:
result = {}
for input_type, pipelines in self.pipelines.items():
result[input_type] = [
{
"name": p.name,
"description": p.description,
"supports_question": p.supports_question,
"estimated_cost": p.estimated_cost_per_call
}
for p in pipelines
]
return result
def get_all_supported_types(self) -> set[str]:
return set(self.pipelines.keys())
Registrar pipelines
from openai import OpenAI
import base64
import fitz
client = OpenAI()
def document_extract_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
config = config or PipelineConfig()
doc = fitz.open(path)
pages = []
for i in range(len(doc)):
text = doc[i].get_text().strip()
if text:
pages.append({"page": i + 1, "text": text})
doc.close()
full_text = "\n\n".join(p["text"] for p in pages)
if question:
response = client.chat.completions.create(
model=config.model,
messages=[
{
"role": "system",
"content": "Responde SOLO con base en el contexto. Si no encuentras la respuesta, dilo. Cita la página."
},
{
"role": "user",
"content": f"Documento:\n{full_text[:8000]}\n\nPregunta: {question}"
}
],
max_tokens=config.max_tokens,
temperature=config.temperature
)
return {
"type": "document_qa",
"answer": response.choices[0].message.content,
"pages_processed": len(pages),
"model": config.model
}
return {
"type": "document_extraction",
"pages": pages,
"total_pages": len(pages),
"total_characters": len(full_text)
}
def image_analyze_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
config = config or PipelineConfig()
with open(path, "rb") as f:
b64 = base64.b64encode(f.read()).decode()
prompt = question or "Analiza esta imagen. Describe qué contiene, extrae cualquier texto visible, y clasifica el tipo de contenido."
response = client.chat.completions.create(
model=config.model,
messages=[{
"role": "user",
"content": [
{"type": "text", "text": prompt},
{"type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{b64}"}}
]
}],
max_tokens=config.max_tokens,
temperature=config.temperature
)
return {
"type": "image_analysis",
"analysis": response.choices[0].message.content,
"model": config.model
}
def audio_transcribe_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
config = config or PipelineConfig()
with open(path, "rb") as f:
transcript = client.audio.transcriptions.create(
model="whisper-1",
file=f,
response_format="text"
)
if question:
response = client.chat.completions.create(
model=config.model,
messages=[
{
"role": "system",
"content": "Responde basándote en la transcripción del audio."
},
{
"role": "user",
"content": f"Transcripción:\n{transcript[:6000]}\n\nPregunta: {question}"
}
],
max_tokens=config.max_tokens,
temperature=config.temperature
)
return {
"type": "audio_qa",
"answer": response.choices[0].message.content,
"transcript": transcript,
"model": config.model
}
return {
"type": "audio_transcription",
"transcript": transcript,
"transcript_length": len(transcript)
}
def video_analyze_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
config = config or PipelineConfig()
import cv2
import os
cap = cv2.VideoCapture(path)
fps = cap.get(cv2.CAP_PROP_FPS)
total_frames = int(cap.get(cv2.CAP_PROP_FRAME_COUNT))
duration = total_frames / fps if fps > 0 else 0
target_frames = min(15, max(5, int(duration / 10)))
interval = max(1, int(total_frames / target_frames))
os.makedirs("/tmp/uc_selector_frames", exist_ok=True)
descriptions = []
frame_id = 0
captured = 0
while True:
ret, frame = cap.read()
if not ret:
break
if frame_id % interval == 0 and captured < target_frames:
frame_path = f"/tmp/uc_selector_frames/frame_{frame_id}.jpg"
cv2.imwrite(frame_path, frame)
with open(frame_path, "rb") as f:
b64 = base64.b64encode(f.read()).decode()
ts = frame_id / fps
minutes = int(ts // 60)
seconds = int(ts % 60)
response = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{
"role": "user",
"content": [
{"type": "text", "text": "Describe esta escena en 1-2 oraciones."},
{"type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{b64}"}}
]
}],
max_tokens=100
)
descriptions.append(f"{minutes}:{seconds:02d} - {response.choices[0].message.content}")
captured += 1
frame_id += 1
cap.release()
timeline = "\n".join(descriptions)
summary_prompt = f"Análisis de frames del video:\n\n{timeline}\n\n"
if question:
summary_prompt += f"Pregunta: {question}"
else:
summary_prompt += "Genera un resumen del video con los puntos clave."
response = client.chat.completions.create(
model=config.model,
messages=[{"role": "user", "content": summary_prompt}],
max_tokens=config.max_tokens
)
return {
"type": "video_analysis",
"summary": response.choices[0].message.content,
"frames_analyzed": captured,
"duration_seconds": round(duration, 2),
"model": config.model
}
def text_analyze_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
config = config or PipelineConfig()
with open(path, "r", encoding="utf-8") as f:
text = f.read()
prompt = question or "Analiza este texto. Resume los puntos clave y extrae información relevante."
response = client.chat.completions.create(
model=config.model,
messages=[
{
"role": "user",
"content": f"Texto:\n{text[:8000]}\n\n{prompt}"
}
],
max_tokens=config.max_tokens,
temperature=config.temperature
)
return {
"type": "text_analysis",
"result": response.choices[0].message.content,
"text_length": len(text),
"model": config.model
}
def create_default_registry() -> PipelineRegistry:
registry = PipelineRegistry()
registry.register(Pipeline(
name="document_extractor",
input_type="document",
handler=document_extract_handler,
description="Extrae texto de PDFs y responde preguntas sobre documentos",
supports_question=True,
estimated_cost_per_call=0.01
))
registry.register(Pipeline(
name="image_analyzer",
input_type="image",
handler=image_analyze_handler,
description="Analiza imágenes: describe contenido, extrae texto, clasifica",
supports_question=True,
estimated_cost_per_call=0.005
))
registry.register(Pipeline(
name="audio_transcriber",
input_type="audio",
handler=audio_transcribe_handler,
description="Transcribe audio y responde preguntas sobre el contenido",
supports_question=True,
estimated_cost_per_call=0.02
))
registry.register(Pipeline(
name="video_analyzer",
input_type="video",
handler=video_analyze_handler,
description="Analiza video extrayendo frames y generando resumen",
supports_question=True,
estimated_cost_per_call=0.10
))
registry.register(Pipeline(
name="text_analyzer",
input_type="text",
handler=text_analyze_handler,
description="Analiza archivos de texto y responde preguntas",
supports_question=True,
estimated_cost_per_call=0.005
))
return registry
Paso 4: Router
import time
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("use_case_selector")
class UseCaseRouter:
def __init__(
self,
config: UseCaseSelectorConfig = None,
registry: PipelineRegistry = None,
cost_tracker = None
):
self.config = config or UseCaseSelectorConfig()
self.detector = TypeDetector(self.config)
self.registry = registry or create_default_registry()
self.cost_tracker = cost_tracker
self.history: list[dict] = []
def route(self, path: str, question: str = None) -> dict:
start_time = time.time()
validation = self.detector.validate(path)
if not validation["valid"]:
return {
"status": "error",
"error": validation["message"],
"detection": validation["detection"]
}
detection = validation["detection"]
input_type = detection["type"]
has_question = question is not None and question.strip() != ""
pipeline = self.registry.get_pipeline(input_type, has_question)
if not pipeline:
return {
"status": "error",
"error": f"No hay pipeline registrado para tipo '{input_type}'",
"detection": detection
}
logger.info(f"Routing: {path} → {pipeline.name} (type={input_type}, question={has_question})")
try:
result = pipeline.handler(
path,
question=question,
config=self.config.default_pipeline_config
)
duration_ms = round((time.time() - start_time) * 1000, 2)
entry = {
"status": "success",
"detection": detection,
"pipeline": pipeline.name,
"result": result,
"duration_ms": duration_ms,
"estimated_cost": pipeline.estimated_cost_per_call,
"question": question
}
if self.cost_tracker:
self.cost_tracker.track_chat(
self.config.default_pipeline_config.model,
input_tokens=1000,
output_tokens=500
)
self.history.append(entry)
logger.info(f"Success: {pipeline.name} in {duration_ms}ms")
return entry
except Exception as e:
duration_ms = round((time.time() - start_time) * 1000, 2)
error_entry = {
"status": "error",
"detection": detection,
"pipeline": pipeline.name,
"error": str(e),
"error_type": type(e).__name__,
"duration_ms": duration_ms
}
self.history.append(error_entry)
logger.error(f"Error in {pipeline.name}: {e}")
if self.config.fallback_enabled:
return self._try_fallback(path, question, error_entry)
return error_entry
def _try_fallback(self, path: str, question: str, original_error: dict) -> dict:
logger.info("Attempting fallback...")
try:
with open(path, "rb") as f:
content = f.read()
if len(content) < 10000:
text_content = content.decode("utf-8", errors="replace")
prompt = f"Analiza este contenido:\n{text_content[:6000]}"
if question:
prompt += f"\n\nPregunta: {question}"
response = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
max_tokens=500
)
return {
"status": "fallback",
"original_error": original_error,
"result": {
"type": "fallback_analysis",
"analysis": response.choices[0].message.content
}
}
except Exception:
pass
return original_error
def route_multiple(self, paths: list[str], question: str = None) -> list[dict]:
return [self.route(path, question) for path in paths]
def get_stats(self) -> dict:
total = len(self.history)
successes = sum(1 for h in self.history if h["status"] == "success")
errors = sum(1 for h in self.history if h["status"] == "error")
fallbacks = sum(1 for h in self.history if h["status"] == "fallback")
durations = [h["duration_ms"] for h in self.history if "duration_ms" in h]
total_cost = sum(h.get("estimated_cost", 0) for h in self.history)
pipeline_counts = {}
for h in self.history:
p = h.get("pipeline", "unknown")
pipeline_counts[p] = pipeline_counts.get(p, 0) + 1
return {
"total_requests": total,
"successes": successes,
"errors": errors,
"fallbacks": fallbacks,
"success_rate": round(successes / total * 100, 1) if total > 0 else 0,
"avg_duration_ms": round(sum(durations) / len(durations), 2) if durations else 0,
"total_estimated_cost": round(total_cost, 4),
"by_pipeline": pipeline_counts
}
Paso 5: Cost Tracker Integrado
class UseCaseCostTracker:
def __init__(self):
self.records: list[dict] = []
def track(self, pipeline_name: str, input_type: str, cost: float, duration_ms: float):
self.records.append({
"pipeline": pipeline_name,
"input_type": input_type,
"cost": cost,
"duration_ms": duration_ms,
"timestamp": time.time()
})
def summary(self) -> dict:
total_cost = sum(r["cost"] for r in self.records)
by_pipeline = {}
for r in self.records:
p = r["pipeline"]
if p not in by_pipeline:
by_pipeline[p] = {"cost": 0, "count": 0}
by_pipeline[p]["cost"] += r["cost"]
by_pipeline[p]["count"] += 1
return {
"total_cost": round(total_cost, 4),
"total_requests": len(self.records),
"by_pipeline": {
k: {"cost": round(v["cost"], 4), "count": v["count"]}
for k, v in by_pipeline.items()
}
}
def check_budget(self, max_daily_cost: float) -> dict:
today_records = [
r for r in self.records
if r["timestamp"] > time.time() - 86400
]
today_cost = sum(r["cost"] for r in today_records)
return {
"today_cost": round(today_cost, 4),
"daily_budget": max_daily_cost,
"remaining": round(max_daily_cost - today_cost, 4),
"within_budget": today_cost < max_daily_cost,
"usage_percent": round(today_cost / max_daily_cost * 100, 1) if max_daily_cost > 0 else 0
}
Paso 6: Main — Orquestación Completa
class UseCaseSelector:
def __init__(self, config: UseCaseSelectorConfig = None):
self.config = config or UseCaseSelectorConfig()
self.registry = create_default_registry()
self.cost_tracker = UseCaseCostTracker()
self.router = UseCaseRouter(
config=self.config,
registry=self.registry,
cost_tracker=self.cost_tracker
)
def process(self, path: str, question: str = None) -> dict:
return self.router.route(path, question)
def process_batch(self, paths: list[str], question: str = None) -> list[dict]:
return self.router.route_multiple(paths, question)
def info(self) -> dict:
return {
"supported_types": list(self.config.supported_types),
"available_pipelines": self.registry.list_pipelines(),
"config": {
"model": self.config.default_pipeline_config.model,
"max_file_size_mb": self.config.max_file_size_mb,
"fallback_enabled": self.config.fallback_enabled,
"cache_enabled": self.config.default_pipeline_config.enable_cache
}
}
def stats(self) -> dict:
return {
"routing": self.router.get_stats(),
"costs": self.cost_tracker.summary()
}
def add_pipeline(self, pipeline: Pipeline) -> None:
self.registry.register(pipeline)
Uso completo
selector = UseCaseSelector()
print(selector.info())
result = selector.process("contrato.pdf")
print(f"Tipo: {result['detection']['type']}")
print(f"Pipeline: {result['pipeline']}")
print(f"Resultado: {result['result']}")
result = selector.process(
"contrato.pdf",
question="¿Cuál es la cláusula de penalización?"
)
print(f"Respuesta: {result['result']['answer']}")
result = selector.process("producto.jpg")
print(f"Análisis: {result['result']['analysis']}")
result = selector.process("reunion.mp3")
print(f"Transcripción: {result['result']['transcript'][:200]}")
result = selector.process("clase.mp4")
print(f"Resumen: {result['result']['summary']}")
print(selector.stats())
Extensión 1: Auto-Pipeline Selection con LLM
En lugar de solo usar la extensión del archivo, usar un LLM para decidir qué pipeline aplicar basándose en el contenido:
def smart_pipeline_selection(
path: str,
question: str,
detection: dict,
available_pipelines: dict
) -> str:
pipelines_desc = "\n".join(
f"- {p['name']}: {p['description']}"
for type_pipelines in available_pipelines.values()
for p in type_pipelines
)
response = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{
"role": "user",
"content": (
f"Archivo: {path}\n"
f"Tipo detectado: {detection['type']}\n"
f"Extensión: {detection['extension']}\n"
f"Tamaño: {detection['size_mb']}MB\n"
f"Pregunta del usuario: {question or 'Ninguna'}\n\n"
f"Pipelines disponibles:\n{pipelines_desc}\n\n"
"¿Qué pipeline es el más apropiado? Responde SOLO con el nombre del pipeline."
)
}],
max_tokens=30,
temperature=0
)
return response.choices[0].message.content.strip()
Extensión 2: Batch Routing con Prioridades
Para procesar múltiples archivos con diferentes prioridades:
from dataclasses import dataclass
from typing import Optional
@dataclass
class RoutingRequest:
path: str
question: Optional[str] = None
priority: int = 5
class BatchRouter:
def __init__(self, selector: UseCaseSelector):
self.selector = selector
self.queue: list[RoutingRequest] = []
def add(self, path: str, question: str = None, priority: int = 5):
self.queue.append(RoutingRequest(path=path, question=question, priority=priority))
def process_all(self) -> list[dict]:
sorted_queue = sorted(self.queue, key=lambda r: r.priority)
results = []
for request in sorted_queue:
result = self.selector.process(request.path, request.question)
result["priority"] = request.priority
results.append(result)
self.queue.clear()
return results
def process_by_type(self) -> dict[str, list[dict]]:
by_type: dict[str, list[RoutingRequest]] = {}
for req in self.queue:
detection = self.selector.router.detector.detect(req.path)
input_type = detection["type"]
if input_type not in by_type:
by_type[input_type] = []
by_type[input_type].append(req)
results = {}
for input_type, requests in by_type.items():
results[input_type] = []
for req in sorted(requests, key=lambda r: r.priority):
result = self.selector.process(req.path, req.question)
results[input_type].append(result)
self.queue.clear()
return results
Uso:
batch = BatchRouter(selector)
batch.add("urgente.pdf", priority=1)
batch.add("producto1.jpg", priority=5)
batch.add("producto2.jpg", priority=5)
batch.add("reunion.mp3", question="¿Qué decisiones se tomaron?", priority=3)
results = batch.process_all()
for r in results:
print(f"[P{r['priority']}] {r['pipeline']}: {r['status']}")
Troubleshooting
Problema 1: Tipo de archivo no detectado
Síntoma: TypeDetector.detect retorna "unknown" para un archivo válido.
Solución: Agregar la extensión al config:
config = UseCaseSelectorConfig()
config.type_extensions[".xlsx"] = "document"
config.type_extensions[".pptx"] = "document"
Problema 2: Pipeline incorrecto seleccionado
Síntoma: Un PDF escaneado se enruta a document_extractor pero no hay texto que extraer.
Solución: Agregar detección de contenido después del tipo:
def detect_pdf_subtype(path: str) -> str:
import fitz
doc = fitz.open(path)
total_text = sum(len(doc[i].get_text().strip()) for i in range(len(doc)))
total_images = sum(len(doc[i].get_images()) for i in range(len(doc)))
doc.close()
if total_text < 100 and total_images > 0:
return "scanned_document"
return "text_document"
Problema 3: Videos largos consumen demasiado
Síntoma: Un video de 2 horas genera 100+ llamadas a Vision API.
Solución: El video_analyze_handler ya limita frames con target_frames = min(15, ...). Para más control:
config = PipelineConfig()
config.max_tokens = 300
selector = UseCaseSelector(UseCaseSelectorConfig(
default_pipeline_config=config,
max_file_size_mb=50.0
))
Problema 4: Costo acumulado sin control
Síntoma: El cost tracker muestra gasto alto pero no hay alertas.
Solución:
budget = selector.cost_tracker.check_budget(max_daily_cost=5.0)
if not budget["within_budget"]:
print(f"ALERTA: Presupuesto excedido. Gastado: ${budget['today_cost']}")
Checklist de Completitud
Funcionalidad core
- Detecta PDF y enruta a document pipeline
- Detecta imagen y enruta a image pipeline
- Detecta audio y enruta a audio pipeline
- Detecta video y enruta a video pipeline
- Detecta texto y enruta a text pipeline
- Maneja archivos no soportados con error claro
- Soporta pregunta opcional para Q&A
Robustez
- Valida que el archivo existe
- Verifica tamaño máximo
- Maneja errores de API con fallback
- Registra historial de operaciones
- Trackea costos estimados
Extensibilidad
- Se pueden registrar pipelines nuevos
- La configuración es flexible
- Soporta batch processing
Ejercicios
Ejercicio 1: Pipeline personalizado
Agrega un pipeline para archivos CSV que lea el archivo, analice las columnas con un LLM, y genere un resumen estadístico.
Ver solución
import csv
def csv_analyze_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
config = config or PipelineConfig()
with open(path, "r", encoding="utf-8") as f:
reader = csv.reader(f)
headers = next(reader)
rows = list(reader)
sample = rows[:20]
sample_text = "\n".join([",".join(headers)] + [",".join(row) for row in sample])
prompt = (
f"Archivo CSV con {len(rows)} filas y {len(headers)} columnas.\n"
f"Columnas: {', '.join(headers)}\n\n"
f"Muestra de datos:\n{sample_text}\n\n"
)
if question:
prompt += f"Pregunta: {question}"
else:
prompt += "Analiza estos datos: describe las columnas, identifica patrones, y genera un resumen."
response = client.chat.completions.create(
model=config.model,
messages=[{"role": "user", "content": prompt}],
max_tokens=config.max_tokens
)
return {
"type": "csv_analysis",
"result": response.choices[0].message.content,
"rows": len(rows),
"columns": headers
}
selector.add_pipeline(Pipeline(
name="csv_analyzer",
input_type="text",
handler=csv_analyze_handler,
description="Analiza archivos CSV",
supports_question=True,
estimated_cost_per_call=0.005
))
Ejercicio 2: Multi-file routing
Extiende el selector para que reciba múltiples archivos y los combine en un solo análisis cuando sean del mismo tema.
Ver solución
def process_related_files(
selector: UseCaseSelector,
paths: list[str],
question: str = None
) -> dict:
individual_results = []
for path in paths:
result = selector.process(path)
if result["status"] == "success":
individual_results.append({
"path": path,
"type": result["detection"]["type"],
"summary": str(result["result"])[:500]
})
if len(individual_results) > 1:
summaries = "\n\n".join(
f"[{r['type']}] {r['path']}:\n{r['summary']}"
for r in individual_results
)
prompt = f"Análisis de {len(individual_results)} archivos:\n\n{summaries}\n\n"
if question:
prompt += f"Pregunta: {question}"
else:
prompt += "Genera un análisis integrado que conecte la información de todos los archivos."
response = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
max_tokens=600
)
return {
"type": "multi_file_analysis",
"integrated_result": response.choices[0].message.content,
"individual_results": individual_results,
"files_processed": len(individual_results)
}
return {
"type": "single_file",
"results": individual_results
}
Resumen
El Use Case Selector es:
- Un router que detecta tipo de input y selecciona el pipeline correcto
- Un registry que gestiona pipelines disponibles y permite agregar nuevos
- Un tracker que monitorea costos y rendimiento
- La base del Document Analyzer del Módulo 8
Componentes implementados:
| Componente | Clase | Responsabilidad |
|---|---|---|
| Config | UseCaseSelectorConfig, PipelineConfig | Configuración flexible |
| Detector | TypeDetector | Detectar tipo de archivo |
| Registry | PipelineRegistry | Gestionar pipelines |
| Router | UseCaseRouter | Seleccionar y ejecutar |
| Cost | UseCaseCostTracker | Rastrear costos |
| Main | UseCaseSelector | Orquestar todo |
Próximo módulo: Módulo 8 — Proyecto Final: Document Analyzer Multimodal. El Use Case Selector que construiste aquí se integra como componente de routing en un sistema completo que procesa documentos, extrae datos con vision, indexa para Q&A con RAG, y genera audio con TTS.