⚙️

ETL y Big Data - Glue, EMR, Kinesis

Glue, EMR, Data Pipeline y procesamiento de datos

⏱️ Tiempo estimado de lectura: 25 minutos

AWS Glue - ETL Serverless

AWS Glue es un servicio de ETL (Extract, Transform, Load) serverless para preparar y transformar datos para análisis.

Componentes principales:

Glue Data Catalog:
- Repositorio central de metadata
- Almacena schemas, tablas, particiones, ubicaciones
- Usado por Athena, Redshift Spectrum, EMR
- Versionado automático de schemas

Glue Crawlers:
- Escanea fuentes de datos automáticamente
- Infiere schemas y crea tablas en Data Catalog
- Detecta particiones y actualiza metadata
- Soporta S3, RDS, DynamoDB, JDBC sources

Glue ETL Jobs:
- Código Spark o Python generado automáticamente
- Transforma datos entre formatos
- Escala automáticamente según volumen de datos
- Pago por segundo de DPU (Data Processing Unit) usado

Características:
- Serverless: No infraestructura que gestionar
- Auto-scaling: Ajusta workers automáticamente
- Job bookmarks: Procesa solo datos nuevos
- Desarrollo: Visual ETL designer o code editor

Puntos Clave

  • Glue Data Catalog es la metadata store central de AWS
  • Crawlers ejecutan on-demand o scheduled
  • ETL jobs pueden ser Spark (Scala/Python) o Python Shell
  • Glue Studio ofrece interfaz visual drag-and-drop
  • Job bookmarks previenen re-procesamiento de datos

💻 Job de Glue para transformar datos

# Glue ETL Job - Convertir CSV a Parquet
import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job

args = getResolvedOptions(sys.argv, ['JOB_NAME'])
sc = SparkContext()
glueContext = GlueContext(sc)
job = Job(glueContext)
job.init(args['JOB_NAME'], args)

# Leer datos desde Data Catalog
datasource = glueContext.create_dynamic_frame.from_catalog(
    database = "mi_database",
    table_name = "ventas_csv"
)

# Transformar: filtrar y agregar columna
transformed = datasource.filter(lambda x: x["cantidad"] > 0)

# Escribir como Parquet particionado
glueContext.write_dynamic_frame.from_options(
    frame = transformed,
    connection_type = "s3",
    connection_options = {"path": "s3://mi-bucket/ventas-parquet/"},
    format = "parquet",
    format_options = {"compression": "snappy"},
    transformation_ctx = "write_parquet"
)

job.commit()

AWS Glue DataBrew

Glue DataBrew es un servicio visual de preparación de datos sin código para limpiar y normalizar datos.

Características:
- Visual: Interfaz point-and-click sin código
- 250+ transformaciones: Limpieza, normalización, filtrado
- Data profiling: Análisis automático de calidad de datos
- Recipes: Secuencias de transformaciones reutilizables
- Preview: Ver resultados antes de aplicar a todo el dataset

Casos de uso:
- Limpiar datos antes de análisis o ML
- Normalizar formatos (fechas, nombres, direcciones)
- Detectar y manejar valores faltantes
- Identificar outliers y anomalías
- Preparar datos sin escribir código

Diferencias con Glue ETL:
- DataBrew: Visual, no-code, para analistas
- Glue ETL: Code-based, más flexible, para ingenieros
- DataBrew genera Recipes que pueden automatizarse
- Ambos pueden integrarse en pipelines

Puntos Clave

  • DataBrew ideal para analistas sin experiencia en programación
  • Data profiling identifica problemas de calidad automáticamente
  • Recipes pueden exportarse y reutilizarse
  • Soporta sampling para trabajar con subsets grandes
  • Puede escribir resultados a S3, Redshift, Data Catalog

Amazon EMR - Elastic MapReduce

Amazon EMR es una plataforma de big data para procesar grandes cantidades de datos usando frameworks open-source como Hadoop, Spark, Hive, Presto.

Frameworks soportados:
- Apache Spark: Procesamiento in-memory, ML, streaming
- Apache Hadoop: MapReduce, HDFS
- Apache Hive: SQL sobre Hadoop
- Apache HBase: Base de datos NoSQL distribuida
- Presto: Queries SQL interactivos
- Flink, Phoenix, Zeppelin, Livy

Arquitectura:
- Master Node: Coordina cluster, gestiona YARN ResourceManager
- Core Nodes: Ejecutan tasks, almacenan datos en HDFS
- Task Nodes: Solo ejecutan tasks, sin almacenamiento (opcionales)

Opciones de deployment:
- Transient clusters: Se crean, procesan datos, y se terminan
- Long-running clusters: Permanecen activos para workloads continuos
- EMR Serverless: Sin gestión de clusters, auto-scaling
- EMR on EKS: Ejecuta Spark jobs en Kubernetes

Almacenamiento:
- HDFS: Local en Core Nodes (efímero)
- EMRFS: S3 como storage persistente (recomendado)
- Instance Store: Para resultados temporales
- EBS: Para HDFS de Core Nodes

Puntos Clave

  • EMR más económico que ejecutar Spark/Hadoop self-managed
  • Usar S3 (EMRFS) para storage persistente, no HDFS
  • Spot Instances ideales para Task Nodes (sin pérdida de datos)
  • EMR Notebooks para desarrollo interactivo con Spark
  • Auto-scaling puede agregar/quitar Task Nodes según carga

💻 Crear y ejecutar jobs en EMR

# Crear cluster EMR con Spark
aws emr create-cluster \n  --name "Mi-Cluster-Spark" \n  --release-label emr-6.10.0 \n  --applications Name=Spark Name=Hive \n  --instance-type m5.xlarge \n  --instance-count 3 \n  --use-default-roles \n  --log-uri s3://mi-bucket/emr-logs/ \n  --bootstrap-actions Path=s3://mi-bucket/bootstrap.sh

# Agregar step (job) al cluster
aws emr add-steps \n  --cluster-id j-XXXXXXXXXXXXX \n  --steps Type=Spark,Name="Procesamiento de datos",\nActionOnFailure=CONTINUE,\nArgs=[--deploy-mode,cluster,--master,yarn,\ns3://mi-bucket/spark-job.py,s3://mi-bucket/input/,\ns3://mi-bucket/output/]

Amazon Kinesis Data Streams

Kinesis Data Streams permite capturar, procesar y almacenar streams de datos en tiempo real.

Conceptos clave:
- Stream: Conjunto de shards que procesan datos
- Shard: Unidad de capacidad (1MB/s entrada, 2MB/s salida)
- Producer: Envía records al stream (SDK, Agent, Firehose)
- Consumer: Lee records del stream (Lambda, EC2, Kinesis Analytics)
- Record: Datos + Partition Key + Sequence Number

Capacidad:
- Provisioned mode: Especificas número de shards manualmente
- On-demand mode: Escala automáticamente según throughput
- 1 shard: 1MB/s in, 2MB/s out, 1000 records/s

Retención:
- Default: 24 horas
- Extendible: hasta 365 días
- Replay: Re-procesar datos históricos

Partition Key:
- Determina a qué shard va el record
- Mismo partition key = mismo shard (orden garantizado)
- Distribuir keys uniformemente previene hot shards

Casos de uso:
- Real-time analytics (dashboards en vivo)
- Log aggregation
- IoT data ingestion
- Clickstream analysis

Puntos Clave

  • Kinesis Data Streams NO es serverless (provisionar shards)
  • Data persiste (no se borra al consumir) durante retention period
  • Enhanced Fan-Out permite múltiples consumers sin afectar throughput
  • KCL (Kinesis Client Library) simplifica consumo con checkpointing
  • ProvisionedThroughputExceeded = necesitas más shards o mejor partition key

Amazon Kinesis Data Firehose

Kinesis Data Firehose es un servicio totalmente gestionado para cargar datos de streaming a destinos de almacenamiento y analytics.

Características:
- Fully managed: Sin administración de shards o scaling
- Near real-time: Buffer de 60s mínimo o 1MB mínimo
- Auto-scaling: Maneja cualquier throughput
- Transformación: Lambda puede transformar records en tránsito
- Conversión: Puede convertir formatos (JSON a Parquet/ORC)

Destinos:
- AWS: S3, Redshift, OpenSearch, Kinesis Data Analytics
- Third-party: Datadog, Splunk, New Relic, MongoDB
- HTTP Endpoints: Custom destinations

Buffer configuración:
- Buffer size: 1MB - 128MB
- Buffer interval: 60s - 900s
- Firehose espera a que se cumpla size O interval (lo que ocurra primero)

Transformaciones:
- Lambda puede modificar records antes de delivery
- Built-in conversion: JSON → Parquet/ORC
- Compression: GZIP, ZIP, Snappy

Diferencia con Kinesis Data Streams:
- Firehose = Entrega a destinos (load)
- Data Streams = Procesamiento custom en real-time

Puntos Clave

  • Firehose es serverless y totalmente gestionado
  • Near real-time (no real-time verdadero como Data Streams)
  • No almacena datos (solo los entrega al destino)
  • Para Redshift: primero escribe a S3, luego COPY a Redshift
  • Failed records pueden enviarse a S3 backup bucket

AWS Data Pipeline vs Glue vs EMR

AWS ofrece múltiples servicios para procesamiento de datos. Es importante saber cuándo usar cada uno.

AWS Data Pipeline:
- Orquestación de workflows de datos
- Mueve y transforma datos entre servicios AWS
- Basado en instancias EC2 o EMR (no serverless)
- Schedule-based execution
- Casos de uso: Backups regulares, procesamiento batch legacy

AWS Glue:
- ETL serverless
- Ideal para transformaciones simples a medianas
- Data Catalog central
- Auto-scaling, pago por segundo
- Casos de uso: Preparación de data lakes, conversión de formatos

Amazon EMR:
- Big data processing con frameworks open-source
- Requiere gestionar clusters (no serverless)
- Máximo control y flexibilidad
- Casos de uso: ML complejo, procesamiento masivo con Spark/Hadoop

Cuándo usar cada uno:
- Glue: Transformaciones SQL-like, ETL serverless, costo moderado
- EMR: Spark/Hadoop jobs complejos, ML, control total
- Data Pipeline: Orquestación legacy, migraciones on-prem → AWS

Puntos Clave

  • Glue es la opción más simple y serverless para ETL
  • EMR para workloads que requieren Spark/Hadoop completo
  • Data Pipeline es legacy, considerar Step Functions + Glue
  • Para streaming: Kinesis Data Streams o Firehose
  • Para batch: Glue (serverless) o EMR (más potente)