Lab práctico · Semana 9: Datos, analítica y machine learning: ingesta, transformación y servicios de IA

Streaming a un data lake: Kinesis Data Streams, Data Firehose a Parquet y consultas con Athena

⏱ 90-120 minDificultad: mediaTask statements: 3.5

Qué vas a construir

Un pipeline de datos completo, en pequeño, como los que describe el task statement 3.5:

  1. Un Kinesis Data Stream (1 shard, modo provisioned) que recibe ventas en JSON en tiempo real.
  2. Un Amazon Data Firehose que lee el stream, convierte el JSON a Parquet usando el esquema del AWS Glue Data Catalog y lo escribe en S3 particionado por fecha (dt=AAAA-MM-DD).
  3. Una tabla en el catálogo que consultas con Amazon Athena.
  4. Una segunda parte por lotes: un CSV histórico que conviertes a Parquet particionado con un CTAS de Athena, comparando los bytes escaneados antes y después.
flowchart LR
  CS["CloudShell (productor)"] -->|"put-record JSON"| KDS["Kinesis Data Stream (1 shard)"]
  KDS --> FH["Data Firehose: JSON a Parquet"]
  GC["Glue Data Catalog: lab18.ventas"] -.->|"esquema"| FH
  FH --> S3["S3: ventas/dt=.../*.parquet"]
  CSV["S3: CSV histórico"] --> CTAS["Athena CTAS"]
  CTAS --> PQ["S3: historico-parquet/tienda=.../"]
  S3 --> ATH["Athena SQL"]
  PQ --> ATH

Antes de empezar

  • Inicia sesión con tu usuario de IAM Identity Center o un usuario de IAM con MFA (nunca con el usuario raíz) con permisos para S3, Kinesis, Firehose, Glue, Athena e IAM.
  • Abre AWS CloudShell en la región eu-south-2 (Europa, España). Kinesis Data Streams, Data Firehose, Glue y Athena están disponibles allí.
  • Si tu cuenta usa el plan gratuito y la consola te impide crear el stream o el Firehose por no estar incluidos, necesitarás pasar al plan de pago para hacer el lab (consulta el módulo 01). No uses créditos que no quieras gastar.

Prepara las variables (si CloudShell se reinicia, vuelve a ejecutar este bloque). set +H desactiva la expansión del historial de bash para que el carácter ! de los prefijos de Firehose no dé problemas.

set +H
export AWS_REGION=eu-south-2
export CUENTA=$(aws sts get-caller-identity --query Account --output text)
export BUCKET=lab18-datalake-$CUENTA
export STREAM=lab18-ventas
export FH=lab18-ventas-parquet
export ROL=lab18-firehose-rol
echo "Cuenta $CUENTA, bucket $BUCKET, región $AWS_REGION"

Paso 1: el bucket del data lake

Qué y por qué: S3 es el almacenamiento del data lake. Los buckets nuevos ya tienen Block Public Access activado y cifrado SSE-S3 por defecto.

aws s3api create-bucket --bucket $BUCKET \
  --create-bucket-configuration LocationConstraint=$AWS_REGION

aws s3api get-public-access-block --bucket $BUCKET

Comprueba que las cuatro opciones de Block Public Access aparecen a true.

Paso 2: el Kinesis Data Stream

Qué y por qué: un stream con 1 shard admite hasta 1 MB/s o 1.000 registros/s de escritura, de sobra para el lab. Usamos modo provisioned para ver el concepto de shard (en producción con tráfico impredecible elegirías on-demand).

aws kinesis create-stream --stream-name $STREAM --shard-count 1 \
  --stream-mode-details StreamMode=PROVISIONED

aws kinesis wait stream-exists --stream-name $STREAM

export STREAM_ARN=$(aws kinesis describe-stream-summary --stream-name $STREAM \
  --query StreamDescriptionSummary.StreamARN --output text)

aws kinesis describe-stream-summary --stream-name $STREAM \
  --query 'StreamDescriptionSummary.[StreamStatus,OpenShardCount,RetentionPeriodHours]'

Deberías ver ACTIVE, 1 shard y 24 horas de retención (el valor por defecto; ampliable hasta 365 días con coste).

Paso 3: base de datos y tabla en el Glue Data Catalog (con Athena)

Qué y por qué: Firehose necesita el esquema de una tabla del Glue Data Catalog para convertir JSON a Parquet, y Athena usará esa misma tabla para consultar. Crearemos la base de datos y la tabla con sentencias DDL de Athena, que las registra en el catálogo.

Primero, una función de ayuda que lanza una consulta, espera a que termine y muestra el estado, los bytes escaneados y el resultado:

consulta() {
  local id estado
  id=$(aws athena start-query-execution --query-string "$1" \
    --result-configuration OutputLocation=s3://$BUCKET/athena-resultados/ \
    --query QueryExecutionId --output text)
  while true; do
    estado=$(aws athena get-query-execution --query-execution-id $id \
      --query QueryExecution.Status.State --output text)
    case $estado in SUCCEEDED|FAILED|CANCELLED) break ;; esac
    sleep 2
  done
  echo "Estado: $estado"
  aws athena get-query-execution --query-execution-id $id \
    --query 'QueryExecution.[Status.StateChangeReason,Statistics.DataScannedInBytes]' --output text
  if [ "$estado" = "SUCCEEDED" ]; then
    aws athena get-query-results --query-execution-id $id --output text \
      --query 'ResultSet.Rows[*].Data[*].VarCharValue'
  fi
}

Crea la base de datos y la tabla de destino, en Parquet y particionada por dt:

consulta "CREATE DATABASE IF NOT EXISTS lab18"

consulta "CREATE EXTERNAL TABLE lab18.ventas (
  id string,
  tienda string,
  producto string,
  unidades int,
  importe double,
  ts string)
PARTITIONED BY (dt string)
STORED AS PARQUET
LOCATION 's3://$BUCKET/ventas/'"

aws glue get-table --database-name lab18 --name ventas \
  --query 'Table.StorageDescriptor.Columns[*].[Name,Type]' --output table

El último comando demuestra que la tabla vive en el Glue Data Catalog, el metastore compartido por Athena, Firehose, EMR y Redshift Spectrum.

Paso 4: el rol de IAM para Firehose

Qué y por qué: Firehose asume un rol para leer del stream, leer el esquema de Glue y escribir en S3. Aplicamos mínimo privilegio: solo este stream, esta tabla y este bucket.

cat > confianza-firehose.json <<'EOF'
{
  "Version": "2012-10-17",
  "Statement": [{
    "Effect": "Allow",
    "Principal": { "Service": "firehose.amazonaws.com" },
    "Action": "sts:AssumeRole"
  }]
}
EOF

cat > permisos-firehose.json <<EOF
{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Sid": "LeerStream",
      "Effect": "Allow",
      "Action": ["kinesis:DescribeStream", "kinesis:DescribeStreamSummary",
                 "kinesis:GetShardIterator", "kinesis:GetRecords", "kinesis:ListShards"],
      "Resource": "$STREAM_ARN"
    },
    {
      "Sid": "LeerEsquema",
      "Effect": "Allow",
      "Action": ["glue:GetTable", "glue:GetTableVersion", "glue:GetTableVersions"],
      "Resource": [
        "arn:aws:glue:$AWS_REGION:$CUENTA:catalog",
        "arn:aws:glue:$AWS_REGION:$CUENTA:database/lab18",
        "arn:aws:glue:$AWS_REGION:$CUENTA:table/lab18/ventas"
      ]
    },
    {
      "Sid": "EscribirS3",
      "Effect": "Allow",
      "Action": ["s3:AbortMultipartUpload", "s3:GetBucketLocation", "s3:GetObject",
                 "s3:ListBucket", "s3:ListBucketMultipartUploads", "s3:PutObject"],
      "Resource": ["arn:aws:s3:::$BUCKET", "arn:aws:s3:::$BUCKET/*"]
    }
  ]
}
EOF

aws iam create-role --role-name $ROL \
  --assume-role-policy-document file://confianza-firehose.json
aws iam put-role-policy --role-name $ROL --policy-name lab18-permisos \
  --policy-document file://permisos-firehose.json

export ROL_ARN=$(aws iam get-role --role-name $ROL --query Role.Arn --output text)
sleep 15

La pausa final da tiempo a que el rol se propague en IAM antes de usarlo.

Paso 5: el Firehose con conversión a Parquet

Qué y por qué: configuramos el destino S3 con:

  • Prefijo con particionado por fecha estilo Hive (dt=AAAA-MM-DD), que Athena entiende como partición. Al usar expresiones en el prefijo, Firehose exige un ErrorOutputPrefix.
  • Buffering de 64 MB o 60 s (con conversión de formato, el tamaño mínimo es 64 MiB; en el lab se cumplirá antes el intervalo).
  • Conversión: deserializador OpenX JSON, esquema de lab18.ventas y serializador Parquet con Snappy.
cat > destino.json <<EOF
{
  "RoleARN": "$ROL_ARN",
  "BucketARN": "arn:aws:s3:::$BUCKET",
  "Prefix": "ventas/dt=!{timestamp:yyyy-MM-dd}/",
  "ErrorOutputPrefix": "errores/!{firehose:error-output-type}/dt=!{timestamp:yyyy-MM-dd}/",
  "BufferingHints": { "SizeInMBs": 64, "IntervalInSeconds": 60 },
  "CompressionFormat": "UNCOMPRESSED",
  "DataFormatConversionConfiguration": {
    "Enabled": true,
    "SchemaConfiguration": {
      "RoleARN": "$ROL_ARN",
      "DatabaseName": "lab18",
      "TableName": "ventas",
      "Region": "$AWS_REGION"
    },
    "InputFormatConfiguration": { "Deserializer": { "OpenXJsonSerDe": {} } },
    "OutputFormatConfiguration": { "Serializer": { "ParquetSerDe": { "Compression": "SNAPPY" } } }
  }
}
EOF

aws firehose create-delivery-stream \
  --delivery-stream-name $FH \
  --delivery-stream-type KinesisStreamAsSource \
  --kinesis-stream-source-configuration KinesisStreamARN=$STREAM_ARN,RoleARN=$ROL_ARN \
  --extended-s3-destination-configuration file://destino.json

La compresión global queda en UNCOMPRESSED porque la compresión la aplica el serializador Parquet (Snappy). Espera a que el Firehose esté activo (uno o dos minutos):

until [ "$(aws firehose describe-delivery-stream --delivery-stream-name $FH \
  --query DeliveryStreamDescription.DeliveryStreamStatus --output text)" = "ACTIVE" ]; do
  echo "Creando..."; sleep 10
done
echo "Firehose ACTIVE"

Si el estado pasa a CREATING_FAILED, revisa el rol (paso 4) y vuelve a crearlo tras borrarlo con aws firehose delete-delivery-stream --delivery-stream-name $FH.

Paso 6: producir eventos en tiempo real

Qué y por qué: simulamos cinco tiendas enviando ventas. La partition key es la tienda: todas las ventas de una tienda caen en el mismo shard y conservan el orden. Con la AWS CLI v2, --cli-binary-format raw-in-base64-out permite pasar el JSON tal cual (sin codificarlo en base64).

productos=(portatil raton teclado monitor auriculares)
for i in $(seq 1 100); do
  t=$(( RANDOM % 5 + 1 ))
  p=${productos[$(( RANDOM % 5 ))]}
  u=$(( RANDOM % 4 + 1 ))
  imp=$(( u * (RANDOM % 200 + 10) ))
  aws kinesis put-record --stream-name $STREAM \
    --partition-key "tienda-$t" \
    --cli-binary-format raw-in-base64-out \
    --data "{\"id\":\"v$i\",\"tienda\":\"T$t\",\"producto\":\"$p\",\"unidades\":$u,\"importe\":$imp.5,\"ts\":\"$(date -u +%FT%TZ)\"}" \
    --query SequenceNumber --output text > /dev/null
done
echo "100 ventas enviadas"

Tarda alrededor de un minuto. Mientras tanto, piensa: si fueran 50.000 ventas por segundo, ¿cuántos shards necesitarías en modo provisioned? ¿Y en on-demand?

Paso 7: ver los ficheros Parquet en S3

Espera unos dos minutos (intervalo de buffer de 60 s más el procesamiento) y lista el bucket:

aws s3 ls s3://$BUCKET/ventas/ --recursive --human-readable
aws s3 ls s3://$BUCKET/errores/ --recursive

Deberías ver uno o varios objetos bajo ventas/dt=AAAA-MM-DD/ (fecha en UTC) y ningún objeto en errores/. Si aparecen errores de formato, descarga uno y revisa el mensaje: suele ser un campo JSON que no coincide con el tipo de la tabla.

Paso 8: consultar con Athena

Registra las particiones nuevas y consulta. MSCK REPAIR TABLE descubre los prefijos dt=... que siguen el estilo Hive.

consulta "MSCK REPAIR TABLE lab18.ventas"

consulta "SELECT tienda, count(*) AS ventas, round(sum(importe), 2) AS total
FROM lab18.ventas
GROUP BY tienda
ORDER BY total DESC"

consulta "SELECT count(*) FROM lab18.ventas WHERE dt = '$(date -u +%F)'"

Fíjate en la segunda línea de salida de cada consulta: son los bytes escaneados (lo que factura Athena). En la consola de Athena (editor de consultas, base de datos lab18) puedes repetir las consultas y ver el mismo dato como Data scanned.

Paso 9: por lotes, de CSV a Parquet con CTAS

Qué y por qué: transformación por lotes, el ejemplo literal del temario (.csv to .parquet). Generamos un CSV histórico de 50.000 filas, lo consultamos tal cual y después lo convertimos a Parquet particionado por tienda con un CTAS.

{
  echo "id,tienda,producto,unidades,importe,fecha"
  awk 'BEGIN { srand(42); split("portatil raton teclado monitor auriculares", p, " ");
    for (i = 1; i <= 50000; i++) {
      printf "h%d,T%d,%s,%d,%.2f,2025-%02d-%02d\n", i, int(rand()*5)+1, p[int(rand()*5)+1],
        int(rand()*4)+1, rand()*900+10, int(rand()*12)+1, int(rand()*28)+1 } }'
} > historico.csv
ls -lh historico.csv
aws s3 cp historico.csv s3://$BUCKET/csv/historico/historico.csv

consulta "CREATE EXTERNAL TABLE lab18.historico_csv (
  id string, tienda string, producto string, unidades int, importe double, fecha string)
ROW FORMAT DELIMITED FIELDS TERMINATED BY ','
STORED AS TEXTFILE
LOCATION 's3://$BUCKET/csv/historico/'
TBLPROPERTIES ('skip.header.line.count'='1')"

consulta "SELECT round(sum(importe), 2) FROM lab18.historico_csv WHERE tienda = 'T3'"

Anota los bytes escaneados: con CSV, Athena lee el fichero entero. Ahora el CTAS (las columnas de partición van al final del SELECT):

consulta "CREATE TABLE lab18.historico_parquet
WITH (
  format = 'PARQUET',
  write_compression = 'SNAPPY',
  external_location = 's3://$BUCKET/historico-parquet/',
  partitioned_by = ARRAY['tienda']
) AS
SELECT id, producto, unidades, importe, fecha, tienda
FROM lab18.historico_csv"

aws s3 ls s3://$BUCKET/historico-parquet/ --recursive --human-readable

consulta "SELECT round(sum(importe), 2) FROM lab18.historico_parquet WHERE tienda = 'T3'"

Compara: la misma suma, pero Athena solo lee la partición tienda=T3 y solo la columna importe de ficheros comprimidos. Los bytes escaneados bajan de forma drástica. Esa es la receta de coste de Athena: columnar + compresión + particionado.

Comprueba que funciona

  • aws s3 ls s3://$BUCKET/ventas/ --recursive muestra objetos Parquet bajo dt=....
  • La consulta del paso 8 devuelve cinco tiendas con un total de 100 ventas.
  • La consulta sobre historico_parquet devuelve el mismo total que sobre historico_csv y escanea muchos menos bytes.
  • En la consola de AWS Glue → Data Catalog → Tables aparecen ventas, historico_csv y historico_parquet en la base de datos lab18.

Limpieza

Hazla en este orden (primero lo que cobra por hora). Ejecútala aunque algún paso haya fallado.

aws firehose delete-delivery-stream --delivery-stream-name $FH
aws kinesis delete-stream --stream-name $STREAM --enforce-consumer-deletion

consulta "DROP TABLE IF EXISTS lab18.historico_parquet"
consulta "DROP TABLE IF EXISTS lab18.historico_csv"
consulta "DROP TABLE IF EXISTS lab18.ventas"
consulta "DROP DATABASE IF EXISTS lab18"

aws s3 rm s3://$BUCKET --recursive
aws s3api delete-bucket --bucket $BUCKET

aws iam delete-role-policy --role-name $ROL --policy-name lab18-permisos
aws iam delete-role --role-name $ROL

rm -f confianza-firehose.json permisos-firehose.json destino.json historico.csv

Comprueba que no queda nada:

aws kinesis list-streams --query StreamNames
aws firehose list-delivery-streams --query DeliveryStreamNames
aws glue get-databases --query 'DatabaseList[?Name==`lab18`]'

Las tres listas deben estar vacías (o no contener los recursos del lab). El DROP TABLE de una tabla externa no borra los datos de S3; por eso se vacía el bucket aparte.

Preguntas para pensar como arquitecto

  1. El equipo de fraude quiere procesar cada venta en menos de un segundo y, a la vez, seguir guardándolas en Parquet. ¿Qué cambiarías en este diseño?
Respuesta

Añadir un segundo consumidor al mismo Kinesis Data Stream: una función Lambda (o una aplicación de Managed Service for Apache Flink) que evalúe cada registro en tiempo real. Firehose sigue siendo otro consumidor que archiva en S3. Esa es la ventaja de Kinesis Data Streams frente a Firehose directo: varios consumidores independientes sobre el mismo flujo y posibilidad de relectura. Si hubiera muchos consumidores, activarías enhanced fan-out.

  1. Si solo necesitaras llevar los eventos a S3 en Parquet, sin tiempo real ni relecturas, ¿qué eliminarías para reducir coste y operación?
Respuesta

El Kinesis Data Stream. Los productores escribirían directamente en Firehose (Direct PUT), que no cobra por shard y hora ni requiere dimensionar nada. Es la opción con LEAST operational overhead para «entregar datos en streaming a S3».

  1. Los productores están en subredes privadas sin salida a Internet. ¿Cómo llegan a Kinesis de forma segura y sin NAT Gateway?
Respuesta

Con un interface VPC endpoint (AWS PrivateLink) de Kinesis Data Streams (o de Firehose) en esas subredes, con un security group que permita HTTPS desde los productores y una política de endpoint restrictiva. Los permisos siguen controlándose con roles de IAM. Para S3, un gateway endpoint gratuito.

  1. Varios equipos van a consultar este data lake y el de marketing no debe ver el importe de las ventas. ¿Qué servicio usarías?
Respuesta

AWS Lake Formation: registras la ubicación de S3, y concedes a marketing permisos sobre la tabla excluyendo la columna importe (o con un data filter por filas y columnas). Athena, Redshift Spectrum, EMR y Amazon Quick respetan esos permisos. Con políticas de bucket de S3 no puedes filtrar columnas.

  1. El CTAS ha convertido 50.000 filas. Si cada noche llegaran 500 GB de CSV, ¿seguirías usando Athena CTAS?
Respuesta

Para cargas grandes y recurrentes es más adecuado un job de AWS Glue (Spark serverless, con job bookmarks para procesar solo lo nuevo y orquestado con un trigger o Step Functions) o EMR Serverless si ya hay código Spark. Athena CTAS es cómodo para conversiones puntuales o moderadas, pero tiene límites (por ejemplo, un máximo de particiones por consulta) y no está pensado como motor ETL recurrente.


Volver al módulo