aj
ajteaches
En esta página
TutorialEngineeringDataEngineeringDistributedProgrammingAWSApacheSparkETLCloud

Distributed Pipeline with Apache Spark and AWS Glue

A�
Alejandro Jaimes 🥭22 min read · 16 de junio de 2026 · 09:25 p. m.
Distributed Pipeline with Apache Spark and AWS Glue

Build a complete Bronze → Silver → Gold ETL pipeline over a dataset of one million real transactions, using AWS Glue Studio, Step Functions, and Athena — entirely from the web console, with no infrastructure to manage.

alejojaimes/ajteaches-tutorials

Python · 0

Hola, estimado lector, como prerrequisito se recomienda tener noción básica de Apache Spark y PySpark, por lo que se invita cordialmente a revisar el siguiente artículo Apache Spark: Introduction, ya que este es el core del ejercicio desarrollado durante esta guia.


¿Qué aprenderás?

A diseñar y ejecutar un pipeline ETL distribuido siguiendo la arquitectura medallion (Bronze → Silver → Gold), a justificar y modelar un esquema estrella para analítica, y a orquestar jobs de Spark de forma automática con manejo de errores.

¿Qué construirás?

Un Data Lake funcional en S3 con datos limpios, tipados y modelados en esquema estrella, capaz de responder preguntas de negocio con SQL sobre un millón de filas reales.

¿Cómo lo harás?

Paso a paso, desde la consola web de AWS (Glue Studio, Step Functions y Athena), sin instalar nada adicional ni gestionar infraestructura. ¡Manos a la obra!

Image by the author

Entorno Virtual

Antes de instalar cualquier paquete, se debe aislar el entorno de Python para no contaminar el sistema global. Desde la raíz del proyecto:

bash
python -m venv venv

Activación del entorno:

# Windows (PowerShell)
venv\Scripts\Activate.ps1

# macOS / Linux
source venv/bin/activate

Con el entorno activo (el prompt de la terminal debe mostrar el prefijo (venv)), se instalan las dependencias necesarias para preparar el dataset:

bash
pip install pandas openpyxl
Image by the author

Estructura del repositorio

Antes de abrir la consola de AWS, o realizar alguna otra acción, es importante crearse un repositorio en Github que obedezca la siguiente estructura de carpetas y nombramiento de artefactos, para manejar un estándard.

bash
spark-glue-workshop/
│
├── README.md
├── .gitignore
├── venv/
├── data/
│   ├── online_retail_II.zip
│   ├── online_retail_II.xlsx
│   ├── online_retail_ii_2009_2010.csv
│   └── online_retail_ii_2010_2011.csv
│
├── scripts/
│   ├── prepare_csv.py
│   ├── job_bronze_to_silver.py
│   └── job_silver_to_gold.py
│
├── sql/
│   ├── create_athena_tables.sql
│   └── business_questions.sql
│
└── evidence/
    ├── step_0_budget.png
    ├── step_2_s3_bucket.png
    ├── step_5_glue_silver_succeeded.png
    ├── step_6_glue_gold_succeeded.png
    ├── step_7_step_functions.png
    └── step_8_athena_queries.png

La carpeta evidence almacena las capturas del paso a paso, durante cada ejecución para realizar el archivo readme final. El cual, evidencia el funcionamiento del pipeline.

Archivo .gitignore

El entorno virtual y los archivos pesados del dataset (el .zip descargado, el .xlsx original y los CSV generados) no deben subirse al repositorio: son binarios grandes, no aportan valor como código versionado, y pueden regenerarse en cualquier momento siguiendo los pasos de esta sección. Debe crearse un archivo .gitignore en la raíz del proyecto con el siguiente contenido:

bash
# Entorno virtual
venv/

# Dataset (se descarga y se genera localmente, no se versiona)
data/*.zip
data/*.xlsx
data/*.csv

# Cache de Python
__pycache__/
*.pyc

Las carpetas data/, evidence/ y venv/ deben existir localmente, pero solo evidence/ se sube al repositorio con contenido. data/ debe quedar vacía en Git (puede agregarse un .gitkeep si se desea conservar la carpeta en el repositorio).

Una vez creado el repositorio y realizada esta estructura (sin las imágenes) inicializar el repositorio con el siguienre comando.

bash
git init
git add .
git commit -m "chore: initial repository structure for spark-glue workshop"

El README como bitácora

El archivo README.md no se escribe todo al final: se va completando paso a paso, en paralelo con cada captura de pantalla. Esto cumple dos propósitos: deja registro de en qué momento se tomó cada evidencia, y obliga a documentar el trabajo a medida que se avanza, en vez de reconstruir la memoria al final.

El contenido base, antes de empezar con el desarrollo del taller debe ser:

# Spark / AWS Glue Workshop

Pipeline ETL distribuido Bronze → Silver → Gold sobre el dataset Online Retail II,
construido con AWS Glue Studio, Step Functions y Athena.

## Evidencia del taller

A partir de aquí, cada vez que el tutorial indique 📸 y se guarde una captura en evidence/, debe agregarse al final del README una entrada con el formato:

### <Nombre del paso>
![<descripción breve>](evidence/<nombre_archivo>.png)

Por ejemplo, al completar el Paso 0 más adelante, se añade al README:

### Paso 0 — Alerta de presupuesto
![Budget creado en Billing](evidence/step_0_budget.png)

Y así sucesivamente con cada captura indicada en el resto del tutorial. Al llegar al PASO 9, el README debe tener seis secciones de evidencia, una por cada captura tomada (PASO 0, 2, 5, 6, 7 y 8).

git add README.md evidence/step_0_budget.png
git commit -m "docs: add budget alert evidence"

A partir de este punto, cada commit de evidencia en el tutorial incluye también README.md en el git add, ya que ambos archivos se actualizan juntos.


Dataset: Online Retail II

Se trabajará con Online Retail II, un dataset real de transacciones de una tienda del Reino Unido, con cerca de 1.067.000 filas correspondientes a dos años (2009–2011).

Se debe contar con una cuenta de Kaggle, para poder descargar el dataset.

Kaggle | Online Retail Dataset

Descarga: https://www.kaggle.com/datasets/lakshmi25npathi/online-retail-dataset

El archivo se descarga como un .zip (online_retail_II.zip). Debe descomprimirse dentro de la carpeta data/, dejando el archivo online_retail_II.xlsx en esa misma ubicación

# Dentro de la carpeta data/
unzip online_retail_II.zip

El Excel contiene dos hojas, una por cada año del dataset:

Hoja en Excel

Filas Aprox

Year 2009-2010

525.461

Year 2009-2010

541.910

Cada hoja debe convertirse en un CSV independiente, con un nombre relacionado a su hoja de origen, en minúsculas, sin espacios y en formato snake_case. Este script debe guardarse como scripts/prepare_csv.py y ejecutarse localmente (con el entorno virtual activo):

python
# scripts/prepare_csv.py
import pandas as pd

SHEET_TO_FILENAME = {
    "Year 2009-2010": "online_retail_ii_2009_2010.csv",
    "Year 2010-2011": "online_retail_ii_2010_2011.csv",
}

xls = pd.ExcelFile("data/online_retail_II.xlsx")

for sheet_name, filename in SHEET_TO_FILENAME.items():
    df = xls.parse(sheet_name)
    output_path = f"data/{filename}"
    df.to_csv(output_path, index=False)
    print(f"{filename} generated: {len(df):,} rows")

Luego ejecutar desde la terminal con python scripts/prepare_csv.py

Image by the author

Al finalizar, la carpeta data/ debe contener dos archivos: online_retail_ii_2009_2010.csv y online_retail_ii_2010_2011.csv. Ambos se suben por separado a la capa Bronze en el PASO 3, y Glue los leerá como un único DataFrame particionado automáticamente al apuntar a la carpeta completa.

Image by the author

Esquema del dataset (idéntico en ambos archivos):

Columna

Tipo

Descripción

Invoice

string

Nº de factura. Si empieza por C es una devolución

StockCode

string

Código de producto

Description

string

Nombre del producto

Quantity

int

Unidades (negativo en devoluciones)

InvoiceDate

timestamp

Fecha y hora

Price

double

Precio unitario en £

Customer ID

double

ID de cliente (puede ser nulo)

Country

string

País del cliente


Arquitectura del Pipeline

Image by the author

Todo el flujo es orquestado por Step Functions.

Se implementa un modelo LakeHouse el cual se establece en una arquitectura Medallion de tres stages, Bronze -> Silver -> Gold. Durante cada etapa se evidenciará el enriquecimiento de los datos, hasta llegar a gold con la implementación del modelo para hacer consultas analíticas.


Con todo esto, vamos manos a la obra con nuestra cuenta de AWS, configurando un presupuesto.

Step 0 - Alerta de presupuesto

  1. Iniciar sesión en https://console.aws.amazon.com.
  2. En la barra de búsqueda superior, escribir Billing and Cost Management y hacer clic en el resultado.
  3. En el menú lateral izquierdo, hacer clic en Budgets.
  4. Hacer clic en el botón naranja Create budget.
  5. En la sección Budget setup, seleccionar Use a template (simplified).
  6. Dentro de Templates, seleccionar Monthly cost budget.
  7. Completar los siguientes campos:

    • Budget name: spark-glue-workshop-budget
    • Enter your budgeted amount ($): 2.00
    • Email recipients: dirección de correo electrónico
  8. Hacer clic en Create budget.

El budget quedará listado con su estado. AWS enviará un correo automáticamente cuando el gasto real alcance el 85% y cuando la proyección llegue al 100%.

⚠️ Esto debe completarse antes de crear cualquier recurso. AWS cobra por uso real, y olvidar limpiar al final puede generar sorpresas en la factura. Un presupuesto de cinco dólares es más que suficiente para todo el taller.

📸 Captura del budget recién creado → evidence/step_0_budget.png

bash
git add evidence/step_0_budget.png
git commit -m "docs: add budget alert evidence"

Step 1 - Tags

Para que todos los recursos queden relacionados entre sí y sea posible rastrear el costo del taller, se utiliza un prefijo común y un set de tags en cada recurso que se cree.

  • Prefijo del proyecto: spark-glue-workshop
  • Región: debe elegirse una sola y usarse de forma consistente (recomendado: us-east-1). Antes de cada paso, conviene verificar en la esquina superior derecha de la consola que la región seleccionada es la correcta.
  • Nombre del bucket: debe ser único en todo AWS. Se recomienda el formato spark-glue-workshop-datalake-<iniciales>-01, por ejemplo: spark-glue-workshop-datalake-jac-01.

Tags que deben añadirse manualmente en cada recurso:

Key

Value

Project

spark-glue-workshop

Environment

dev

Owner

<nombre>

Course

parallel-distributed-computing

ManagedBy

workshop

En Billing → Cost Explorer los costos pueden agruparse por Tag: Project para ver exactamente cuánto costó el uso de esos servicios. Es la forma de hacer tracking de los servicios utilizados y saber su costo.


Step 2 - Crear DataLake (Bucket S3)

El bucket S3 es la base de todo el pipeline. En él viven el CSV crudo, los datos procesados y los resultados de Athena.

  1. En la barra de búsqueda superior, escribir S3 y hacer clic en el servicio.
  2. Hacer clic en el botón naranja Create bucket.
  3. En el campo Bucket name, escribir: spark-glue-workshop-datalake-<iniciales>-01
  4. En AWS Region, seleccionar la misma región elegida en el STEP 1.
  5. En la sección Block Public Access settings, dejar todos los checkboxes marcados (es la configuración recomendada, no debe modificarse).
  6. En la sección Bucket Versioning, seleccionar Enable.
  7. Bajar hasta la sección Tags y hacer clic en Add tag. Añadir las 5 tags de la tabla del PASO 1, una por una.
  8. Hacer clic en Create bucket.

Crear las carpetas del Data Lake

Una vez creado el bucket, hacer clic en su nombre para entrar. Deben crearse cinco carpetas repitiendo estos pasos para cada una:

  1. Hacer clic en el botón Create folder.
  2. Escribir el nombre de la carpeta.
  3. Hacer clic en Create folder.

Las cinco carpetas necesarias:

Folder

Para qué sirve

bronze/

CSV crudos generados apartir del excel descargado de Kaggle.

silver/

Parquet limpio y tipado

gold/

Modelo estrella (fact + dimensiones)

athena-results/

Resultados de consultas de Athena

temp/

Carpeta temporal interna de Glue

📸 Captura del bucket con las 5 carpetas visibles → evidence/step_2_s3_bucket.png

Agregar al README la entrada correspondiente:

### Paso 2 — Bucket S3 con las 5 carpetas
![Bucket del Data Lake](evidence/step_2_s3_bucket.png)
git add README.md evidence/step_2_s3_bucket.png
git commit -m "docs: add S3 bucket creation evidence"

Step 3 - Subir datasets ( Bronze -> Lawyer)

  1. Dentro del bucket, hacer clic en la carpeta bronze/ para entrar.
  2. Hacer clic en el botón Upload.
  3. Hacer clic en Add files.
  4. Seleccionar ambos archivos generados en el paso anterior desde el computador local: online_retail_ii_2009_2010.csv y online_retail_ii_2010_2011.csv.
  5. Hacer clic en el botón naranja Upload al final de la página.
  6. Esperar a que la barra de progreso llegue al 100%. Los archivos son grandes y puede tardar algunos minutos según la conexión.
  7. Al finalizar, ambos archivos quedarán listados dentro de bronze/ con su tamaño y fecha.

💡 Como los dos archivos quedan en la misma carpeta bronze/, Glue los leerá como un único DataFrame al apuntar al STEP 5 a la carpeta completa (s3://<bucket>/bronze/), particionado automáticamente. No es necesario unirlos manualmente


Step 4 - IAM Rol para Glue

Glue necesita un rol con permisos para leer y escribir en el bucket. Sin esto, los jobs fallarán con AccessDenied.

  1. En la barra de búsqueda superior, escribir IAM y hacer clic en el servicio.
  2. En el menú lateral izquierdo, hacer clic en Roles.
  3. Hacer clic en el botón naranja Create role.
  4. En Trusted entity type, seleccionar AWS service.
  5. En la sección Use case, hacer clic en el campo de búsqueda, escribir Glue y seleccionar la opción Glue que aparece.
  6. Hacer clic en Next.
  7. En la barra de búsqueda de permisos, buscar y marcar el checkbox de AWSGlueServiceRole. Luego buscar y marcar también AmazonS3FullAccess.
  8. Hacer clic en Next.
  9. En Role name, escribir: spark-glue-workshop-role
  10. Bajar hasta la sección Tags y añadir las 5 tags del STEP 1.
  11. Hacer clic en Create role.

Tener presente el nombre de este rol, ya que se seleccionará al configurar cada Job de Glue.


Step 5 - Glue Jobs: Bronze -> Silver

Image by Author

Este es el primer Job que se realizará, el cual lee los CSV en crudo, los limpia, tipifica y lo guarda como parquet. Aquí Spark realiza su trabajo real sobre el clúster.

5.1 Crear el job

  1. En la barra de búsqueda, escribir AWS Glue y hacer clic en el servicio.
  2. En el menú lateral izquierdo, hacer clic en ETL jobs.
  3. Hacer clic en Script editor.
  4. En el popup que aparece, seleccionar Engine: Spark y Start fresh.
  5. Hacer clic en Create script.
  6. En la parte superior, hacer clic sobre el nombre por defecto del job y cambiarlo a: spark-glue-workshop-bronze-to-silver
Image by author
Image by author
Image by author

5.2 Asociar el rol IAM y configurar el job

  1. Hacer clic en la pestaña Job details (en la parte superior, junto a "Script").
  2. Configurar los siguientes campos:

    • IAM Role: buscar y seleccionar spark-glue-workshop-role
    • Glue version: seleccionar Glue 4.0
    • Worker type: seleccionar G.1X
    • Requested number of workers: 2
    • Temporary path: s3://<bucket>/temp/
  3. Bajar hasta la sección Tags y añadir las 5 tags del STEP 1.
  4. Bajar hasta la sección Job parameters y hacer clic en Add new parameter:

    • Key: --BUCKET
    • Value: el nombre del bucket (sin s3://, solo el nombre)

5.3 Pegar el código del job

  1. Regresar a la pestaña Script.
  2. Borrar todo el contenido del editor y pegar lo que se encuentra en el script completo job_bronze_to_silver.py

Importante: logger viene de glue_context.get_logger(), no de un logger genérico de Python. Eso hace que cada línea quede automáticamente en CloudWatch Logs, visible desde la pestaña Runs del job → Logs sin configurar nada adicional. Cada etapa numerada (STEP 1 a STEP 7) corresponde a una transformación lógica del pipeline, y cada try/except evita que un fallo silencioso pase desapercibido: si algo falla, el log indica exactamente en qué paso ocurrió antes de relanzar la excepción (raise) y detener el job.

5.4 Guardar y ejecutar

  1. Hacer clic en Save (esquina superior derecha).
  2. Hacer clic en Run.
  3. Hacer clic en la pestaña Runs para seguir el estado. El job tarda entre 3 y 5 minutos. Esperar hasta ver el estado Succeeded.
  4. Ir a S3 → el bucket → la carpeta silver/. Deberían verse carpetas como year=2010/ y year=2011/, cada una con archivos .parquet.

🧪 Validación teórica/práctica: regresar a Job details, cambiar el número de workers a 4, guardar y ejecutar de nuevo. Comparar los tiempos en la pestaña Runs. Más workers significa más particiones procesadas en paralelo; en algún punto, agregar workers deja de reducir el tiempo proporcionalmente — esa es la Ley de Amdahl en la práctica.

5.5 Validar el logger en CloudWatch

El logger usado en el script no escribe en el mismo lugar que un print(). Glue separa los logs en grupos distintos:

Origen

Log group en CloudWatch

print(...)

/aws-glue/jobs/output

logger.info(...), logger.error(...) (vía glue_context.get_logger())

/aws-glue/jobs/logs-v2

Por eso, revisar Output Logs nunca mostrará las líneas STEP 1, STEP 2, etc.: hay que ir al log group correcto.

  1. En la pestaña Runs del job, ubicar el run que se quiere inspeccionar.
  2. En esa misma fila, hacer clic en el enlace bajo la columna Logs (no en "Output logs"; ese enlace lleva al log group /aws-glue/jobs/output, que solo contiene la salida de Spark y de print()). El enlace de Logs lleva al log group /aws-glue/jobs/logs-v2, donde vive el logger nativo.
  3. Dentro del log group, el log stream tiene como nombre el Job Run ID (un identificador largo que empieza con jr_), no el nombre del job. Hacer clic en ese stream.
  4. Usar el cuadro de Filter events (arriba de la lista de líneas) y escribir STEP. CloudWatch filtrará automáticamente y solo se verán las líneas propias del script:

📸 Captura del run con estado Succeeded → evidence/step_5_glue_silver_succeeded.png

Agregar al README la entrada correspondiente:

markdown
### Step 5 — Job Bronze → Silver completado
![Job Bronze a Silver con estado Succeeded](evidence/step_5_glue_silver_succeeded.png)
bash
git add README.md evidence/step_5_glue_silver_succeeded.png
git commit -m "docs: add glue bronze-to-silver job evidence"

Step 6 — Glue Jobs: Silver → Gold

Este job transforma los datos limpios de Silver en un modelo estrella: una tabla de hechos y tres dimensiones.

What is a Star Schema (and why it's important) - Adam Gilmore

¿Por qué un modelo estrella?

Los datos en Silver ya están limpios y tipados, pero siguen siendo una tabla plana: cada fila repite el nombre del producto, el país del cliente y el detalle de la fecha. Esa repetición funciona para almacenar, pero no para consultar: cualquier agregación por producto, cliente o periodo termina escaneando y agrupando texto repetido en millones de filas.

El modelo estrella separa esa información en dos tipos de tabla:

  • Dimensiones (dim_product, dim_customer, dim_date): catálogos pequeños con los atributos descriptivos (nombre de producto, país, día de la semana) y una clave subrogada (*_sk) generada por Spark.
  • Tabla de hechos (fact_sales): una fila por línea de venta, pero con solo las claves subrogadas y las métricas numéricas (quantity, unit_price, total_amount), sin texto repetido.

Esto trae varias ventajas concretas para este pipeline:

  • Tablas más livianas: fact_sales ocupa una fracción del espacio de Silver, porque el texto repetido (descripciones, países) vive una sola vez en cada dimensión.
  • Joins más baratos: las dimensiones son pequeñas (miles de filas, no millones), así que pueden enviarse completas a cada nodo (broadcast) en vez de mezclar (shuffle) la tabla de hechos completa contra ellas.
  • Consultas más simples y rápidas: preguntas de negocio como "ingresos por producto" o "ventas por país" se resuelven con un JOIN + GROUP BY directo sobre claves numéricas, que es justamente lo que Athena necesita para escanear menos datos en Parquet.
  • Reutilización: las mismas tres dimensiones sirven para responder preguntas distintas (por producto, por cliente, por tiempo) sin reprocesar Silver cada vez.

Las consultas frecuentes que este modelo resuelve de forma directa son: ranking de productos por ingresos, evolución de ventas por mes o año, comparación de ingresos por país, e identificación de los clientes con mayor gasto — todas ellas implementadas más adelante en el STEP 8.

Image by author

6.1 Clonar el job y reemplazar lo necesario

Este job se configura igual que el del STEP 5: mismo motor, mismo rol, mismo tipo de worker. Solo cambia el nombre del job y el código.

  1. En Glue → ETL jobs, hacer clic en el job spark-glue-workshop-bronze-to-silver y dar clic en Clone job.
  2. En la parte superior, hacer clic sobre el nombre del job y cambiarlo a: spark-glue-workshop-silver-to-gold
  3. Hacer clic en la pestaña Job details y repetir exactamente la configuración del PASO 5.2 (rol spark-glue-workshop-role, Glue 4.0, G.1X, 2 workers, mismo Temporary path, las 5 tags y el parámetro --BUCKET).

6.2 Pegar el código del job

  1. Regresar a la pestaña Script.
  2. Borrar todo el contenido del editor y pegar lo que se encuentra en el script completo job_silver_to_gold.py
Image by author
Image by author

Importante: generar las claves subrogadas y hacer los joins provoca shuffle. Pero como las dimensiones son pequeñas (pocos miles de filas), se envuelven en F.broadcast(): Spark las copia a todos los nodos y mueve solo los datos de la tabla de hechos. Esta es la optimización de join más común en pipelines distribuidos.

Image by author

6.3 Guardar y ejecutar

  1. Hacer clic en Save.
  2. Hacer clic en Run.
  1. Ir a la pestaña Runs y esperar a que el estado sea Succeeded (puede tardar 4–6 minutos).
  2. Ir a S3 → el bucket → la carpeta gold/. Deberían verse cuatro carpetas: dim_product/, dim_customer/, dim_date/, fact_sales/.

Image by author
Image by author

💡 Para revisar los logs de este job, el procedimiento es el mismo que en el STEP 5.5: columna Logs de la fila del run (log group /aws-glue/jobs/logs-v2), filtrar por STEP dentro del log stream con nombre jr_....

📸 Captura del run con estado Succeeded → evidence/step_6_glue_gold_succeeded.png

Agregar al README la entrada correspondiente:

markdown
### Paso 6 — Job Silver → Gold completado (modelo estrella)
![Job Silver a Gold con estado Succeeded](evidence/step_6_glue_gold_succeeded.png)
bash
git add README.md evidence/step_6_glue_gold_succeeded.png
git commit -m "docs: add glue silver-to-gold job evidence"

PASO 7 — Orquestación con Step Functions

AWS Step Functions 101 - Mariliis Retter

Hasta ahora los jobs se ejecutaron manualmente y por separado. Ahora se encadenan: primero Bronze→Silver, y cuando termine, Silver→Gold. Si alguno falla, esto se hace visible de inmediato. Eso es un pipeline orquestado.

En vez de armar la state machine arrastrando bloques en el Workflow Studio, se define directamente en código: Step Functions usa un lenguaje declarativo en JSON llamado Amazon States Language (ASL). Definirla como archivo permite versionarla en el repositorio igual que los scripts de Glue, en vez de que solo exista dentro de la consola.

7.1 Agregar el archivo a la estructura del repositorio

Debe crearse la carpeta step-functions/ en la raíz del proyecto, con el archivo state_machine.json dentro:

spark-glue-workshop/
│
├── step-functions/
│   └── state_machine.json
│
├── ...

Contenido de step-functions/state_machine.json:

json
{
  "Comment": "Orquesta el pipeline Bronze -> Silver -> Gold del Spark / AWS Glue Workshop",
  "StartAt": "BronzeToSilver",
  "States": {
    "BronzeToSilver": {
      "Type": "Task",
      "Resource": "arn:aws:states:::glue:startJobRun.sync",
      "Parameters": {
        "JobName": "spark-glue-workshop-bronze-to-silver"
      },
      "Catch": [
        {
          "ErrorEquals": ["States.ALL"],
          "Next": "PipelineFailed"
        }
      ],
      "Next": "SilverToGold"
    },
    "SilverToGold": {
      "Type": "Task",
      "Resource": "arn:aws:states:::glue:startJobRun.sync",
      "Parameters": {
        "JobName": "spark-glue-workshop-silver-to-gold",
        "Arguments": {
          "--BUCKET": "REPLACE_WITH_YOUR_BUCKET_NAME"
        }
      },
      "Catch": [
        {
          "ErrorEquals": ["States.ALL"],
          "Next": "PipelineFailed"
        }
      ],
      "End": true
    },
    "PipelineFailed": {
      "Type": "Fail",
      "Error": "PipelineExecutionFailed",
      "Cause": "Uno de los jobs de Glue fallo. Revisar CloudWatch Logs (/aws-glue/jobs/logs-v2) para el detalle."
    }
  }
}

IMPORTANTE: Antes de usarlo, debe reemplazarse REPLACE_WITH_YOUR_BUCKET_NAME con el nombre real del bucket creado en el STEP 2.

Para pensar: Resource: "arn:aws:states:::glue:startJobRun.sync" es lo que hace que Step Functions espere a que el job de Glue termine antes de avanzar al siguiente estado (el sufijo .sync activa ese comportamiento). Sin él, Step Functions dispararía el job y pasaría inmediatamente al siguiente estado sin esperar el resultado. El bloque Catch en ambos estados captura cualquier error (States.ALL) y redirige al estado PipelineFailed, que detiene la ejecución con un mensaje explicativo en vez de dejarla colgada o fallando en silencio.

bash
git add step-functions/state_machine.json
git commit -m "feat: add step functions state machine definition"

7.2 Crear la state machine a partir del JSON

  1. En la barra de búsqueda, escribir Step Functions y hacer clic en el servicio.
  2. En el menú lateral izquierdo, hacer clic en State machines.
  3. Hacer clic en el botón naranja Create state machine.
  4. Seleccionar Create from Blank (lienzo en blanco).
  5. Colocar el nombre de la step como sf-spark-glue-workshop
  6. Seleccionar el modo Standard y dar clic en continue.
  7. En la parte superior del editor, cambiar del modo Design al modo Code.
  8. Borrar todo el contenido del editor de código y pegar el contenido completo del archivo step-functions/state_machine.json (ya con el nombre del bucket reemplazado).
  9. El editor mostrará automáticamente el diagrama del flujo a la derecha (BronzeToSilverSilverToGold, con la rama de error hacia PipelineFailed), generado a partir del JSON pegado, sin necesidad de armarlo a mano.

7.3 Guardar y crear

  1. Hacer clic en Next (esquina superior derecha).
  2. Completar los campos:

    • State machine name: spark-glue-workshop-orchestrator
    • Permissions: seleccionar Create new role (Step Functions crea automáticamente un rol con los permisos para invocar Glue)
  3. Bajar hasta Tags y añadir las 5 tags del STEP 1.
  4. Hacer clic en Create state machine.

7.4 Ejecutar el pipeline completo

  1. En la página de la state machine, hacer clic en Start execution.
  2. En el popup de input, dejar {} y hacer clic en Start execution.
  3. Observar el grafo en vivo: BronzeToSilver aparecerá en azul (ejecutando). Cuando termine y se ponga en verde, arrancará automáticamente SilverToGold. Si algo falla, el flujo salta a PipelineFailed y el estado correspondiente se pone en rojo; puede hacerse clic sobre él para ver los logs.

Felicitaciones, construiste tu primer pipeline distribuido completo ejecutándose como un diagrama, sin tocar una línea de infraestructura. Como lo hacen los pros 😎

📸 Captura del grafo con ambos estados en verde → evidence/step_7_step_functions.png

Agregar al README la entrada correspondiente:

markdown
### Paso 7 — Orquestación con Step Functions
![Pipeline orquestado con ambos estados completados](evidence/step_7_step_functions.png)
bash
git add README.md evidence/step_7_step_functions.png
git commit -m "docs: add step functions orchestration evidence"

Step 8 — Queries analíticas con Athena

Image by author

Con los datos en Gold (modelo estrella en Parquet), pueden responderse preguntas de negocio con SQL directamente sobre S3, sin mover datos ni levantar ningún servidor.

8.1 Configurar Athena

  1. En la barra de búsqueda, escribir Athena y hacer clic en el servicio.
  2. Hacer clic en Query editor.
  3. Si es la primera vez que se usa Athena, aparece un aviso en la parte superior. Hacer clic en Settings (esquina superior derecha del editor).
  4. Hacer clic en Manage.
  5. En el campo Location of query result, escribir: s3://<bucket>/athena-results/
  6. Hacer clic en Save.

8.2 Registrar la base de datos y las tablas

Estas queries deben guardarse también en sql/create_athena_tables.sql. En el editor de Athena, copiar y ejecutar primero:

sql
CREATE DATABASE IF NOT EXISTS workshop_gold;

Luego, en el desplegable Database del panel izquierdo, seleccionar workshop_gold. Ahora deben ejecutarse las cuatro tablas externas (reemplazando <BUCKET> con el nombre real del bucket):

sql
CREATE EXTERNAL TABLE IF NOT EXISTS workshop_gold.dim_product (
  stock_code  string,
  description string,
  product_sk  bigint)
STORED AS PARQUET
LOCATION 's3://<BUCKET>/gold/dim_product/';

CREATE EXTERNAL TABLE IF NOT EXISTS workshop_gold.dim_customer (
  customer_id bigint,
  country     string,
  customer_sk bigint)
STORED AS PARQUET
LOCATION 's3://<BUCKET>/gold/dim_customer/';

CREATE EXTERNAL TABLE IF NOT EXISTS workshop_gold.dim_date (
  date     date,
  year     int,
  month    int,
  day      int,
  weekday  string,
  date_sk  int)
STORED AS PARQUET
LOCATION 's3://<BUCKET>/gold/dim_date/';

CREATE EXTERNAL TABLE IF NOT EXISTS workshop_gold.fact_sales (
  invoice      string,
  product_sk   bigint,
  customer_sk  bigint,
  date_sk      int,
  quantity     int,
  unit_price   double,
  total_amount double,
  is_return    boolean)
STORED AS PARQUET
LOCATION 's3://<BUCKET>/gold/fact_sales/';

Importante, que se ejecuten una por una las sentencias. Actualizar con el boton de flechita de la esquina superior izquierda.

8.3 Preguntas de negocio

Estas queries deben guardarse en sql/business_questions.sql. Deben ejecutarse una por una en el editor:

P1 — ¿Cuáles son los 10 productos con más ingresos?

sql
SELECT
  p.description,
  ROUND(SUM(f.total_amount), 2) AS revenue
FROM workshop_gold.fact_sales f
JOIN workshop_gold.dim_product p ON f.product_sk = p.product_sk
WHERE f.is_return = false
GROUP BY p.description
ORDER BY revenue DESC
LIMIT 10;

P2 — ¿Cómo evolucionaron las ventas mes a mes?

sql
SELECT
  t.year,
  t.month,
  ROUND(SUM(f.total_amount), 2) AS revenue
FROM workshop_gold.fact_sales f
JOIN workshop_gold.dim_date t ON f.date_sk = t.date_sk
WHERE f.is_return = false
GROUP BY t.year, t.month
ORDER BY t.year, t.month;

P3 — ¿Qué países generan más ingresos?

sql
SELECT
  c.country,
  ROUND(SUM(f.total_amount), 2) AS revenue
FROM workshop_gold.fact_sales f
JOIN workshop_gold.dim_customer c ON f.customer_sk = c.customer_sk
WHERE f.is_return = false
GROUP BY c.country
ORDER BY revenue DESC
LIMIT 15;

P4 — ¿Quiénes son los 5 mejores clientes?

sql
SELECT
  c.customer_id,
  c.country,
  COUNT(DISTINCT f.invoice)      AS num_invoices,
  ROUND(SUM(f.total_amount), 2) AS total_spend
FROM workshop_gold.fact_sales f
JOIN workshop_gold.dim_customer c ON f.customer_sk = c.customer_sk
WHERE f.is_return = false
  AND c.customer_id IS NOT NULL
GROUP BY c.customer_id, c.country
ORDER BY total_spend DESC
LIMIT 5;

📸 Captura de una de las consultas con resultados → evidence/step_8_athena_queries.png

Agregar al README la entrada correspondiente, y así cerrar la sección de evidencia con las seis capturas del taller:

markdown

markdown
### Paso 8 — Consultas analíticas en Athena
![Consulta de negocio ejecutada sobre el modelo Gold](evidence/step_8_athena_queries.png)

bash

bash
git add README.md sql/ evidence/step_8_athena_queries.png
git commit -m "docs: add athena queries and evidence"

Step 9 — Limpieza de recursos

Esto debe hacerse al terminar. S3 y Glue siguen cobrando aunque no se use nada si los recursos quedan activos.

El orden importa:

  1. Glue jobs: En Glue → ETL jobs → seleccionar los dos jobs con el checkbox → hacer clic en ActionsDelete → confirmar.

  2. Step Functions: Entrar a la state machine → hacer clic en Delete (esquina superior derecha) → confirmar escribiendo el nombre.

  3. S3: Entrar al bucket → seleccionar todos los objetos con el checkbox superior → hacer clic en Delete → confirmar escribiendo permanently delete. Con el bucket vacío, regresar a la lista de buckets → seleccionarlo → Delete → confirmar escribiendo el nombre del bucket.

  4. IAM: En IAM → Roles → buscar spark-glue-workshop-role → seleccionarlo → Delete → confirmar. También debe buscarse el rol que creó Step Functions automáticamente (suele llamarse StepFunctions-...) y borrarlo también.

  5. Budget: Puede dejarse (no tiene costo) o eliminarse desde Billing → Budgets → seleccionarlo → Delete.

Troubleshooting

  • AccessDenied en Glue → S3: verificar que el job usa spark-glue-workshop-role y que la URL del bucket en el código es la correcta.
  • Job Silver→Gold falla con error de BUCKET: confirmar que se agregó el parámetro --BUCKET en Job details → Job parameters, sin el prefijo s3://.
  • El job tarda demasiado: subir el número de workers en Job details. Verificar que el destino es Parquet y no CSV.
  • Athena no devuelve filas: revisar que la ruta LOCATION en cada CREATE TABLE apunta exactamente a la carpeta de esa tabla dentro de gold/, incluyendo la barra al final.
  • El nombre del bucket ya existe: los nombres de bucket en S3 son globales en todo AWS. Debe cambiarse el número o las iniciales del sufijo.
  • No llegan correos del budget: los avisos solo se disparan al cruzar el umbral. Si el gasto es menor al 85%, no llega ningún correo.

Glosario

Término

Significado

DPU / Worker

Unidad de cómputo de Glue. G.1X = 4 vCPU, 16 GB RAM

Partición

Trozo del dataset procesado en paralelo por un executor

Shuffle

Movimiento de datos entre nodos. Ocurre en groupBy y join. Costoso.

Broadcast join

Copia una tabla pequeña a todos los nodos para evitar shuffle

Parquet

Formato columnar comprimido. Mucho más rápido que CSV para analítica

Medallion

Bronze (crudo) → Silver (limpio) → Gold (modelado)

Esquema estrella

1 tabla de hechos + N dimensiones

SK (surrogate key)

Clave subrogada generada por Spark, no la clave del origen

DAG

Plan de ejecución que construye Spark antes de correr

Lazy evaluation

Las transformaciones no se ejecutan hasta que llega una acción

Si leyó hasta aquí, gracias — de verdad.

Escribir esto tomó tiempo y saber que alguien lo lee hasta el final lo hace valer la pena. Si algo no quedó claro, tiene dudas, o simplemente quiere debatir algún punto, queda la cajita de los comentarios. Leo todo (casi siempre jajaja).

Si quiere estar al tanto, de cuando salen tutoriales como este — suscríbase. Sin spam, solo contenido valioso cuando haya algo que valga la pena compartir.

Nos vemos en la próxima, espero que haya sido de mucha utilidad.

— Alejo🥭

Comments (1)

  • JM
    Julio Martínezlast month

    Nice post, that's good how can explain this teacher!

Sign in to comment