Sesión 3 — Un data lake sobre almacenamiento S3-compatible¶

En la sesión anterior trabajamos con ficheros Parquet almacenados en HDFS. Ahora conservaremos el mismo formato y el mismo conjunto de datos (TPC-DS SF1), pero cambiaremos el sistema de almacenamiento: los objetos se guardarán en un servicio compatible con la API de Amazon S3.

El objetivo es aprender qué cambia al pasar de un sistema de ficheros distribuido a un almacenamiento de objetos, y qué partes del código pueden mantenerse al migrar posteriormente a Amazon S3 real. La sesión construye un data lake sobre S3: conserva los Parquet como objetos bajo prefijos. Las sesiones 6 y 7 añadirán la gestión de tablas, esquemas y snapshots para construir un lakehouse.

Objetivos¶

Al terminar la sesión deberías poder:

  • explicar la diferencia entre un bloque HDFS y un objeto S3;
  • distinguir un bucket, una clave de objeto y un prefijo;
  • arrancar y detener un almacenamiento S3 local sin modificar el clúster Hadoop;
  • utilizar la API estándar de AWS para listar buckets, consultar objetos y leer una parte de un objeto;
  • copiar los Parquet deterministas de TPC-DS SF1 desde HDFS hasta S3;
  • leer los mismos Parquet desde S3 con PyArrow, Polars y DuckDB;
  • parametrizar un programa para usar RustFS local o Amazon S3 real;
  • borrar y regenerar el data lake sin confundir sus datos con los del warehouse HDFS.

Dónde se ejecuta este notebook. Igual que en la sesión 2, Jupyter se ejecuta directamente en namenode (como luser). Las celdas ejecutables acceden a HDFS, WebHDFS y RustFS por la red del laboratorio. Las órdenes que necesitan el Docker del host (arrancar o detener contenedores, make, docker compose) se muestran como texto para ejecutarlas en una terminal de tu equipo, no como celdas: el kernel de este notebook no tiene el socket de Docker.

Paquetes de Python de esta sesión. Esta sesión necesita boto3, aiobotocore, duckdb, fsspec, polars, pyarrow, s3fs. Respecto a S2 añade boto3, aiobotocore, s3fs. Se instalan en el kernel de namenode en la sección «Instalar las dependencias de la sesión en este kernel»: una celda %%writefile requirements.txt crea el fichero dentro del contenedor —el requirements.txt de tu copia de s3/ está en tu equipo y namenode no lo ve— y otra ejecuta %pip install -r requirements.txt. La imagen de namenode ya trae estos paquetes en su entorno Python del curso, /opt/tcdm/venv, que es el intérprete del kernel y donde instala %pip; lo normal es que %pip sólo confirme que están; ejecuta esas celdas de todos modos: dejan explícito qué necesita la sesión y reparan el entorno si falta algo.

Diapositivas de la sesión¶

Las diapositivas se generan con el paquete jupyter-notebook-slide. El alias %%diapositiva permite mantener en español los tipos usados en la sesión y produce la misma salida HTML en Jupyter, Colab y la referencia HTML publicada.

En cada sesión encontrarás varios tipos de diapositivas que te ayudarán a seguir el hilo:

  • Sección (titulo): abre la sesión y resume qué vamos a ver y qué haremos.
  • A continuación (avance): anuncia el contenido de la parte siguiente, para saber en qué punto del programa estamos.
  • Recapitulación (resumen): condensa la parte anterior en puntos clave para recordarlos y revisarlos después.
  • Pregunta guía (pregunta): plantea cuestiones que orientan lo que viene a continuación; conviene intentar responderlas antes de seguir.
  • Evaluación de la sesión (evaluacion): recuerda cómo se evalúa el trabajo de esta sesión.
In [ ]:
%pip install -q "jupyter-notebook-slide @ git+https://github.com/dsevilla/jupyter-notebook-slide.git"
%load_ext notebook_slide

import notebook_slide as jnbs

# Colores de 26-27/teoria/tcdm.css, adaptados al tema de las sesiones.
jnbs.configure(
    background="#eaf2f8",
    foreground="#1f2933",
    border="#c9d6e1",
    heading="#0c304d",
    subheading="#0c304d",
    link="#174f7a",
    code_background="#eaf2f8",
    code_foreground="#0c304d",
    quote_background="#ffffff",
    font_family="Atkinson Hyperlegible, Inter, Aptos, Segoe UI, Helvetica, Arial, sans-serif",
    font_url="https://fonts.googleapis.com/css2?family=Atkinson+Hyperlegible:ital,wght@0,400;0,700;1,400;1,700&display=swap",
)
jnbs.register_slide_type("avance", "A continuación", "#174f7a")
jnbs.register_slide_type("resumen", "Recapitulación", "#a54467")
jnbs.register_slide_type("pregunta", "Pregunta guía", "#0c304d")
jnbs.register_slide_type("evaluacion", "Evaluación de la sesión", "#b3701a")
jnbs.register_slide_type(
    "titulo",
    "Sección",
    "#174f7a",
    layout="title",
    background="linear-gradient(135deg, #0c304d, #174f7a 62%, #a54467)",
    foreground="#ffffff",
    border="transparent",
    heading="#ffffff",
    subheading="#dceaf4",
)
jnbs.register_alias("diapositiva")

Sesión 3: un data lake sobre almacenamiento S3-compatible

Cómo vamos a recorrer la sesión

  • Mismo TPC-DS SF1 y mismo Parquet: cambia el almacenamiento, de bloques HDFS a objetos S3
  • Arrancamos RustFS, un servicio local compatible con la API de Amazon S3
  • Copiamos raw/tpcds de HDFS a un bucket con boto3
  • Leemos los mismos objetos con PyArrow, Polars y DuckDB y comparamos los resultados
  • Vemos qué habría que cambiar para usar Amazon S3 real
Evaluación de la sesión

Cómo se evalúa esta sesión

  • Reunión individual breve con el profesor en una sesión posterior
  • Se muestra en vivo, en tu ordenador: el bucket con los Parquet copiados, la comprobación de la API S3 y las tres lecturas de date_dim desde S3
  • Se entrega también una memoria breve (una o dos páginas)
  • No hace falta memorizar órdenes: sí explicar qué cambia entre HDFS y un almacenamiento de objetos
A continuación

De rutas HDFS a claves de objeto

  • hdfs:///datalake/raw/tpcds/date_dim/... pasa a ser s3://tcdm-datalake/raw/tpcds/date_dim/...
  • Bucket, clave y prefijo: los «directorios» son sólo texto dentro de la clave
  • HDFS ofrece un sistema de ficheros con bloques replicados; S3, objetos y operaciones HTTP (GET, PUT, HEAD, listar por prefijo)
  • Los Parquet no cambian al copiarlos; RustFS es el servicio S3 local del laboratorio

Del data lake HDFS al data lake S3¶

En S2 los datos se encontraban en una ruta del espacio de nombres HDFS:

hdfs:///datalake/raw/tpcds/date_dim/part-...

En S3 se guardarán en un bucket y una clave de objeto:

s3://tcdm-datalake/raw/tpcds/date_dim/part-...

Se conserva así la misma arquitectura medallion: raw identifica la capa, tpcds la fuente y date_dim la tabla. Esta sesión solo escribe en raw/: las sesiones 5 a 7 añaden las capas silver/ y gold/, pero lo hacen sobre HDFS, no sobre este bucket. Aquí se observa el patrón de prefijos por capa, no un data lake S3 con las tres capas completas.

La apariencia de directorios es una convención: S3 no necesita crear un directorio date_dim, el texto separado por / forma parte de la clave del objeto. Esto tiene consecuencias prácticas:

  • HDFS ofrece operaciones de sistema de ficheros, permisos y bloques replicados entre DataNodes;
  • S3 ofrece objetos identificados por claves dentro de buckets y operaciones HTTP, como GET, PUT, HEAD y listados por prefijo;
  • una herramienta puede mostrar una jerarquía de carpetas aunque el servicio sólo esté almacenando nombres de objetos;
  • la consistencia, las operaciones de renombrado y la gestión de permisos no deben suponerse idénticas a las de HDFS.

Los ficheros Parquet no cambian al copiarlos: siguen teniendo su esquema, metadatos y firma PAR1, pero ahora el lector los obtiene mediante la API de objetos. El prefijo agrupa lógicamente la colección de objetos Parquet que corresponde a date_dim: no hay ningún directorio físico, solo objetos cuya clave comparte ese mismo texto de prefijo.

El servicio local RustFS¶

Se usa RustFS como almacenamiento S3-compatible local: permite trabajar con clientes estándar de AWS mediante un servicio aislado conectado a la red del laboratorio. El Compose de esta sesión está en entorno/compose-s3.yml y contiene únicamente:

Servicio Función ¿Permanece ejecutándose?
rustfs-s3 API S3 y consola web Sí
bucket-setup-s3 Crea el bucket si todavía no existe No, termina correctamente
s3-client Contenedor con AWS CLI, útil desde una terminal Sí

El servicio usa un volumen Docker (tcdm-26-27-rustfs-s3-data). Detener los contenedores no borra los objetos; down -v, en cambio, elimina también ese volumen y debe reservarse para regenerar completamente el ejercicio.

Las credenciales del Compose (tcdm_student / tcdm_student_secret) son específicas del laboratorio local. Con ellas se comprueban el protocolo S3 y el flujo de datos entre HDFS, RustFS y los lectores de la sesión.

A continuación

Arrancar el clúster, RustFS y Jupyter

  • En una terminal de tu equipo: make -C entorno warehouse-up y make -C entorno s3-up
  • Jupyter se arranca dentro de namenode, como en S2: docker exec, su - luser y jupyter lab
  • RustFS no es un DataNode ni un NodeManager: es otro servicio de la red hadoop-cluster
  • Desde el notebook se llega por rustfs:9000; desde tu equipo, por localhost:9000 (consola en 9001)

Arrancar el clúster y RustFS¶

Esta sesión reutiliza los Parquet que S2 dejó en HDFS. Estas órdenes necesitan el Docker del host, así que se ejecutan desde la raíz de la distribución de sesiones, en una terminal de tu equipo, no en este notebook (si ya lo hiciste en S2 y sigue arrancado, puedes saltar la primera):

Animación: la orden make -C entorno s3-up se teclea en una terminal de tu equipo (host)

cd ~/tcdm-public
git pull                       # trae s3/ y entorno/compose-s3.yml
make -C entorno warehouse-up   # si el clúster de S2 no sigue arrancado
make -C entorno s3-up

Este notebook se ejecuta en el Jupyter de namenode, que no arranca solo. Igual que en S2, el servidor se arranca dentro del contenedor.

Abre otra terminal de tu equipo y entra en el contenedor namenode:

Animación: docker exec -it namenode bash se teclea en tu equipo y abre una terminal dentro del contenedor namenode, donde se ejecutan su - luser y jupyter lab

docker exec -it namenode bash

docker exec -it abre una shell interactiva dentro del contenedor: a partir de aquí esa terminal ya no está en tu equipo, sino en namenode, y el indicador lo muestra (root@namenode). Ahí dentro, cambia al usuario de trabajo y arranca Jupyter Lab:

su - luser
jupyter lab --ip=0.0.0.0 --port=8888 --no-browser --IdentityProvider.token=tcdm

su - luser abre una shell de login de luser, que carga /etc/profile.d/tcdm-python.sh y pone /opt/tcdm/venv/bin al principio del PATH: por eso jupyter es el del entorno Python del curso, que ya trae PySpark y las bibliotecas de las sesiones (command -v jupyter lo muestra). --ip=0.0.0.0 hace que Jupyter escuche en la red del contenedor, y compose-hadoop-cluster.yml publica su puerto 8888 en 127.0.0.1 de tu equipo. --IdentityProvider.token=tcdm fija el token de acceso, así que la dirección de conexión es siempre la misma. Jupyter queda en primer plano: deja esa terminal abierta mientras trabajes. Ctrl+C lo detiene y dos exit te devuelven a tu equipo.

El servidor, el kernel y todas las celdas del notebook se ejecutan así dentro de namenode; en tu equipo sólo quedan Visual Studio Code y el fichero del notebook. entorno/Makefile ofrece el atajo make -C entorno jupyter, que hace lo mismo en una sola orden (docker exec -it namenode su - luser -c 'jupyter lab ...'), pero conviene conocer los pasos anteriores para saber dónde se ejecuta cada cosa.

En Visual Studio Code, conecta este notebook con Select Kernel → Select Another Kernel → Existing Jupyter Server, usando la dirección http://127.0.0.1:8888/lab?token=tcdm (la que empieza por 127.0.0.1) y eligiendo el kernel Python 3 (ipykernel). El token es siempre tcdm.

Si algo no funciona —el notebook no conecta, una celda !hdfs no encuentra la orden, un servicio no responde—, consulta «Solución de problemas» en entorno/README.md.

La red hadoop-cluster conecta los contenedores, pero RustFS no es un DataNode ni un NodeManager: es simplemente otro servicio accesible desde namenode por red. Los puertos publicados en el host son:

Servicio Dirección (desde el host)
API S3 http://localhost:9000
Consola RustFS http://localhost:9001

Desde dentro de namenode (donde corre este notebook), el servicio se alcanza por su nombre DNS interno, rustfs, en el puerto 9000 — no por localhost. La consola web es una ayuda visual; las operaciones reproducibles de esta sesión se hacen con boto3, porque es la misma interfaz que se podrá reutilizar con AWS real.

Comprueba desde este notebook que el clúster Hadoop sigue respondiendo (esto sí es una celda ejecutable, porque sólo habla con HDFS):

In [ ]:
!hdfs dfs -ls -d /datalake/raw/tpcds
A continuación

Copiar raw/tpcds de HDFS al bucket

  • Antes: traer los programas de la sesión, instalar dependencias y fijar endpoint y credenciales con %env
  • Cada fichero fluye de hdfs dfs -cat a upload_fileobj, sin pasar por el disco local
  • Se borra antes el prefijo raw/tpcds/ del bucket: repetir el notebook da el mismo conjunto de objetos
  • HDFS no se modifica; al terminar debe haber 24 tablas y al menos 24 objetos

Traer los programas de la sesión al kernel¶

Este notebook se ejecuta dentro de namenode, igual que en S2, y ese contenedor no monta el directorio s3/ de tu distribución local: los cuatro programas de la sesión (check_s3_api.py, read_s3_pyarrow.py, read_s3_polars.py, read_s3_duckdb.py) todavía no existen en el directorio de trabajo del kernel. requirements.txt se crea más abajo con %%writefile, como en S2. La siguiente celda los descarga del repositorio público dsevilla/tcdm-public (rama 26-27), igual que S2 descarga la instantánea del INE.

In [ ]:
from pathlib import Path
from urllib.request import Request, urlopen

SESSION_FILES: tuple[str, ...] = (
    "check_s3_api.py",
    "read_s3_pyarrow.py",
    "read_s3_polars.py",
    "read_s3_duckdb.py",
)
base_url: str = "https://raw.githubusercontent.com/dsevilla/tcdm-public/26-27/s3"

for filename in SESSION_FILES:
    request: Request = Request(f"{base_url}/{filename}", headers={"User-Agent": "Mozilla/5.0"})
    with urlopen(request, timeout=30) as response:
        Path(filename).write_bytes(response.read())

print(f"Descargados {len(SESSION_FILES)} ficheros de la sesión al directorio de trabajo del kernel")

Instalar las dependencias de la sesión en este kernel¶

A partir de aquí todo se ejecuta desde este mismo notebook. La celda siguiente crea requirements.txt con %%writefile dentro de namenode, y la posterior instala sus dependencias con %pip en el intérprete del kernel. Como en S2, la imagen ya las trae y lo normal es que %pip sólo confirme que están.

In [ ]:
%%writefile requirements.txt
# Dependencias de los ejemplos de S3.
# Se utilizan únicamente APIs S3 estándar, de modo que el endpoint local se
# puede sustituir por un bucket real de Amazon S3 sin cambiar las lecturas.
# s3fs se usa en las lecturas de PyArrow y Polars; DuckDB usa httpfs.
# boto3 y aiobotocore se fijan en versiones que comparten botocore 1.43.56.
boto3==1.43.56
aiobotocore==3.9.0
duckdb
fsspec>=2026.2.0
polars
pyarrow
# Esta versión de s3fs proporciona la consulta de metadatos `modified` que
# utilizan los lectores de la sesión.
s3fs>=2026.2.0
In [ ]:
%pip install -q -r requirements.txt

requirements.txt fija versiones compatibles de boto3, aiobotocore, fsspec y s3fs para que la instalación sea reproducible.

Las credenciales y el endpoint de RustFS se fijan una sola vez para el resto del notebook mediante variables de entorno del kernel (%env las deja disponibles tanto para celdas Python como para las celdas ! que lanzan los programas de esta sesión como subprocesos):

In [ ]:
%env AWS_ENDPOINT_URL=http://rustfs:9000
%env AWS_ACCESS_KEY_ID=tcdm_student
%env AWS_SECRET_ACCESS_KEY=tcdm_student_secret
%env AWS_DEFAULT_REGION=us-east-1
%env AWS_S3_ADDRESSING_STYLE=path

Antes de copiar nada comprobamos que RustFS responde y que el bucket existe. Si esta celda falla con un error de conexión, RustFS no está arrancado: ejecuta make -C entorno s3-up en una terminal de tu equipo y repítela.

In [ ]:
import os

import boto3

buckets: list[str] = [
    bucket["Name"]
    for bucket in boto3.client("s3", endpoint_url=os.environ["AWS_ENDPOINT_URL"]).list_buckets()[
        "Buckets"
    ]
]
print(f"Buckets en RustFS: {buckets}")
assert (
    "tcdm-datalake" in buckets
), "No existe el bucket tcdm-datalake: revisa `make -C entorno s3-up`"

Crear el data lake S3 a partir de HDFS¶

El propio notebook copia ahora los datos. Así, al ejecutar las celdas en orden se parte del estado final de S2 y se obtiene exactamente el estado que necesita el resto de esta sesión. El flujo de cada fichero es equivalente a

hdfs dfs -cat fichero.parquet
        │ flujo binario
        ▼
boto3.upload_fileobj(...) -> s3://tcdm-datalake/raw/tpcds/tabla/fichero

La celda de copia enumera recursivamente con hdfs dfs -ls -R y filtra las líneas que empiezan por - (ficheros, no directorios) excluyendo _SUCCESS:

hdfs dfs -ls -R /datalake/raw/tpcds \
  | awk '$1 ~ /^-/ && $NF !~ /\/_SUCCESS$/ { print $NF }'

Para cada ruta, hdfs dfs -cat entrega un flujo binario a boto3; el fichero no pasa por el disco local ni por el host. Antes de subir se elimina sólo el prefijo raw/tpcds/ del bucket para que repetir el notebook produzca el mismo conjunto de objetos y no conserve fragmentos de una ejecución anterior. HDFS permanece intacto: sigue siendo la fuente raw creada en S2.

Copiar los Parquet desde HDFS con la API S3¶

La copia se implementa aquí, de forma visible. subprocess sólo abre el lector nativo de HDFS; la enumeración, la elección de claves y las llamadas a S3 pertenecen a esta celda. upload_fileobj consume la salida de hdfs dfs -cat como un flujo y evita crear una copia local intermedia.

La salida esperada indica 24 tablas y al menos 24 objetos. Puede haber más de un objeto por tabla si Trino decidió escribir varios fragmentos.

In [ ]:
import os
from subprocess import PIPE, Popen, run
from typing import Any

import boto3

HDFS_ROOT: str = "/datalake/raw/tpcds"
S3_BUCKET: str = "tcdm-datalake"
S3_PREFIX: str = "raw/tpcds"

listing: str = run(
    ["hdfs", "dfs", "-ls", "-R", HDFS_ROOT],
    check=True,
    capture_output=True,
    text=True,
).stdout
hdfs_files: list[str] = sorted(
    fields[-1]
    for line in listing.splitlines()
    if (fields := line.split())
    if fields[0].startswith("-") and not fields[-1].endswith("/_SUCCESS")
)
table_names: set[str] = {path.removeprefix(f"{HDFS_ROOT}/").split("/", 1)[0] for path in hdfs_files}
assert len(table_names) == 24, sorted(table_names)
assert hdfs_files, "S2 no dejó ficheros Parquet en HDFS"

s3: Any = boto3.client("s3", endpoint_url=os.environ["AWS_ENDPOINT_URL"])
old_keys: list[str] = [
    item["Key"]
    for page in s3.get_paginator("list_objects_v2").paginate(
        Bucket=S3_BUCKET, Prefix=f"{S3_PREFIX}/"
    )
    for item in page.get("Contents", [])
]
for start in range(0, len(old_keys), 1000):
    batch: list[str] = old_keys[start : start + 1000]
    s3.delete_objects(
        Bucket=S3_BUCKET,
        Delete={"Objects": [{"Key": key} for key in batch], "Quiet": True},
    )

for position, hdfs_path in enumerate(hdfs_files, start=1):
    relative_path: str = hdfs_path.removeprefix(f"{HDFS_ROOT}/")
    s3_key: str = f"{S3_PREFIX}/{relative_path}"
    process: Popen[bytes] = Popen(["hdfs", "dfs", "-cat", hdfs_path], stdout=PIPE)
    assert process.stdout is not None
    try:
        s3.upload_fileobj(process.stdout, S3_BUCKET, s3_key)
    except Exception:
        process.kill()
        process.wait()
        raise
    finally:
        process.stdout.close()
    return_code: int = process.wait()
    if return_code != 0:
        raise RuntimeError(f"hdfs dfs -cat falló para {hdfs_path}: {return_code}")
    print(f"[{position}/{len(hdfs_files)}] {s3_key}")

uploaded_keys: list[str] = [
    item["Key"]
    for page in s3.get_paginator("list_objects_v2").paginate(
        Bucket=S3_BUCKET, Prefix=f"{S3_PREFIX}/"
    )
    for item in page.get("Contents", [])
]
assert len(uploaded_keys) == len(hdfs_files)
print(f"Copiadas {len(hdfs_files)} piezas Parquet de {len(table_names)} tablas.")

Hay dos formas de ver el resultado sin escribir más Python, ambas opcionales:

  • abre la consola de RustFS en http://localhost:9001, entra con tcdm_student / tcdm_student_secret y navega por el bucket tcdm-datalake hasta raw/tpcds/;

  • en una terminal de tu equipo, lista el prefijo con el cliente de AWS del contenedor s3-client, que ya tiene configurados el endpoint y las credenciales:

    docker exec tcdm-s3-client aws s3 ls s3://tcdm-datalake/raw/tpcds/
    
Recapitulación

El mismo raw, ahora como objetos

  • s3://tcdm-datalake/raw/tpcds/<tabla>/... contiene los mismos Parquet que /datalake/raw/tpcds
  • La copia ha usado sólo la API S3 estándar (boto3), configurada con variables AWS_*
  • HDFS sigue intacto: es la fuente raw creada en S2
Pregunta guía
  • Algunos objetos no terminan en .parquet: ¿cómo puede saber un lector que lo son?
  • ¿Hace falta descargar el objeto completo para comprobarlo?

Comprobar la API de AWS con boto3¶

check_s3_api.py muestra las operaciones básicas de un cliente S3 estándar: list_buckets(), list_objects_v2(), head_object() y get_object(Range="bytes=0-3"). Se ejecuta como los demás scripts de la sesión, con python, sin ningún docker exec por delante — el nombre DNS rustfs ya resuelve porque este kernel vive en la misma red que el contenedor.

In [ ]:
!python check_s3_api.py tcdm-datalake raw/tpcds/date_dim --minimum-objects 1

El programa termina imprimiendo TCDM_S3_API_OK después de comprobar que la respuesta de get_object(Range=...) es exactamente b'PAR1': confirma que el objeto remoto comienza con la firma de Parquet y que el acceso por rangos HTTP funciona, sin haber descargado el objeto completo.

A continuación

Tres lectores para los mismos objetos

  • PyArrow: s3fs adaptado con FSSpecHandler
  • Polars: s3fs abre cada objeto como flujo binario y se concatenan los fragmentos
  • DuckDB: su extensión httpfs habla directamente con la API S3
  • Los tres deben dar la misma firma que en S2: filas, claves no nulas, mínimo, máximo y suma

Explorar con PyArrow¶

read_s3_pyarrow.py utiliza s3fs (que implementa el acceso S3 mediante el SDK de AWS) y adapta ese filesystem a PyArrow mediante FSSpecHandler:

PyArrow Parquet
      │
      ▼
FSSpecHandler
      │
      ▼
s3fs → botocore → API S3

El patrón de búsqueda de objetos termina en /*, no en /*.parquet: algunos escritores de Trino generan ficheros Parquet sin esa extensión, y siguen siendo reconocibles por la firma binaria PAR1. Por eso el programa enumera los objetos del prefijo y excluye explícitamente _SUCCESS.

In [ ]:
!python read_s3_pyarrow.py tcdm-datalake raw/tpcds/date_dim \
  --key-column d_date_sk --expected-rows 73049

El programa lista las claves Parquet, lee todos los fragmentos, imprime el esquema y cinco filas, y calcula filas, claves no nulas, mínimo, máximo y suma de la clave. Esas cinco medidas deben coincidir exactamente con las obtenidas al leer la misma tabla desde HDFS en S2: el almacenamiento cambió, el contenido no.

Explorar con Polars¶

read_s3_polars.py usa el mismo s3fs para localizar los objetos, abre cada fichero como flujo binario y lo entrega a polars.read_parquet(), concatenando los fragmentos en un único DataFrame. Esa concatenación facilita el ejercicio, pero reúne todas las filas en memoria; más adelante se usarán lecturas lazy y Spark para colecciones mayores (sesión 4).

In [ ]:
!python read_s3_polars.py tcdm-datalake raw/tpcds/date_dim \
  --key-column d_date_sk --expected-rows 73049

Explorar con DuckDB¶

DuckDB puede consultar Parquet sin crear una tabla permanente. Su extensión nativa httpfs accede directamente a la API S3. read_s3_duckdb.py carga la extensión, configura un secreto S3 con el endpoint, las credenciales y el estilo de URL, y consulta directamente URLs s3:// — el contenido sigue siendo el mismo Parquet, y el endpoint puede ser RustFS local o Amazon S3.

In [ ]:
!python read_s3_duckdb.py tcdm-datalake raw/tpcds/date_dim \
  --key-column d_date_sk --expected-rows 73049

La consulta conceptual detrás del programa es:

DESCRIBE SELECT * FROM read_parquet([...]);
SELECT * FROM read_parquet([...]) LIMIT 5;
SELECT count(*), min(d_date_sk), max(d_date_sk), sum(d_date_sk)
FROM read_parquet([...]);

Aquí el catálogo no sabe nada de date_dim: DuckDB recibe directamente la lista de objetos Parquet (obtenida con list_objects_v2), y la lectura la realiza httpfs, no s3fs.

Comparar los resultados¶

Las tres herramientas deben producir la misma firma (filas, claves_no_nulas, mínimo, máximo, suma): PyArrow entrega una representación columnar explícita, Polars construye un DataFrame y DuckDB ejecuta un plan SQL, pero las tres terminan leyendo los mismos objetos Parquet. Si alguna difiere, sospecha primero de qué prefijo o bucket se está consultando antes de sospechar del propio dato.

Una lectura selectiva por columnas (posible porque Parquet es columnar) reduce los bytes que deben decodificarse, aunque el objeto completo viva en S3. Esta celda sí es Python nativo del kernel, no un script externo:

In [ ]:
import os

import polars as pl
import s3fs

filesystem: s3fs.S3FileSystem = s3fs.S3FileSystem(
    client_kwargs={"endpoint_url": os.environ["AWS_ENDPOINT_URL"]}
)
date_fragments: list[str] = sorted(filesystem.glob("tcdm-datalake/raw/tpcds/date_dim/*"))
assert date_fragments, "No hay fragmentos de date_dim en S3"
with filesystem.open(date_fragments[0], "rb") as stream:
    sample: pl.DataFrame = pl.read_parquet(stream, columns=["d_date", "d_year"])
print(sample.head())

La celda obtiene el nombre real del primer fragmento: el número y los nombres de los ficheros por tabla no están garantizados, sólo el contenido lógico. Más adelante se combinará esta propiedad con filtros y particiones en Spark, Trino e Iceberg (sesiones 4 a 8).

Recapitulación

Tres caminos, la misma firma

  • get_object(Range="bytes=0-3") devuelve PAR1 sin descargar el objeto entero
  • PyArrow y Polars llegan por s3fs; DuckDB, por httpfs: el resultado numérico coincide
  • Leer sólo algunas columnas reduce los bytes que hay que decodificar
  • No hay catálogo: cada lector recibe la lista de objetos del prefijo
A continuación

Del laboratorio a Amazon S3 real

  • Los cuatro programas se configuran con variables estándar AWS_*, nada propio de RustFS
  • Cambian el endpoint, las credenciales, la región y el estilo de direccionamiento (path → virtual)
  • El código de los programas no cambia

Configuración local y configuración AWS¶

Los cuatro programas de esta sesión leen su configuración mediante variables estándar del SDK de AWS, no mediante nada propio de RustFS:

Variable RustFS local Amazon S3
AWS_ENDPOINT_URL http://rustfs:9000 dentro de Docker normalmente no se define
AWS_ACCESS_KEY_ID tcdm_student credencial IAM del laboratorio o rol
AWS_SECRET_ACCESS_KEY secreto local del Compose secreto temporal o rol IAM
AWS_SESSION_TOKEN no se define token temporal del laboratorio, si existe
AWS_DEFAULT_REGION us-east-1 región del bucket
AWS_S3_ADDRESSING_STYLE path dentro de Docker preferentemente virtual

AWS_S3_ADDRESSING_STYLE no es una variable que boto3 o el CLI de AWS lean de forma nativa: son los cuatro programas de esta sesión los que la leen con os.environ.get(...) y la pasan explícitamente a Config(s3={'addressing_style': ...}). Aquí hace falta path porque tcdm-datalake no es un nombre resoluble por DNS fuera del laboratorio; en Amazon S3 real, virtual —el estilo recomendado hoy por AWS— sí resolvería correctamente tcdm-datalake.s3.<región>.amazonaws.com, siempre que el bucket exista en esa región.

Migrar a AWS real consiste en cambiar los valores de %env de más arriba (eliminando AWS_ENDPOINT_URL, usando credenciales IAM reales y AWS_S3_ADDRESSING_STYLE=virtual) y el nombre del bucket; el código de los cuatro programas no cambia.

Recapitulación

Sesión 3

  • Un objeto S3 es una clave dentro de un bucket; el prefijo agrupa, pero no es un directorio
  • El mismo Parquet se lee igual desde HDFS y desde S3: cambia el acceso, no el contenido
  • Parametrizar endpoint y credenciales permite pasar de RustFS a Amazon S3 sin tocar el código
  • Sigue siendo un data lake de ficheros: sin catálogo, esquemas gestionados ni snapshots
Siguiente sesión: Spark como motor de cómputo distribuido sobre estos mismos datos.

Detener, regenerar y limpiar¶

Estas órdenes necesitan el Docker del host, así que se ejecutan en una terminal de tu equipo, no en este notebook (y así un «ejecutar todo» del notebook tampoco puede borrarte los datos por accidente):

Situación Orden Consecuencia
Detener RustFS conservando sus objetos docker compose -f entorno/compose-s3.yml down Los objetos siguen en el volumen Docker; se recupera con make -C entorno s3-up.
Eliminar el laboratorio S3 por completo (destructivo) docker compose -f entorno/compose-s3.yml down -v --remove-orphans Borra también el volumen de RustFS y sus objetos. No elimina los bloques HDFS del clúster Hadoop. Después hay que ejecutar make -C entorno s3-up y volver a ejecutar el notebook.

Este down -v sólo afecta al Compose de S3 (RustFS): no toca el clúster Hadoop ni el HDFS de la sesión 1, que desde la Sesión 1 vive en volúmenes Docker con nombre propios y solo se borra de forma explícita con make -C entorno clean-hdfs.

Si en vez de las órdenes anteriores usas make -C entorno clean, también se borra el volumen de RustFS de esta sesión (clean hace down -v sobre el Compose de S3), aunque conserva el HDFS por el motivo anterior.

Preguntas para interpretar la experiencia¶

  • ¿Qué información reside en la clave de un objeto S3 que en HDFS estaría repartida entre el espacio de nombres del NameNode y los bloques de los DataNodes?
  • ¿Por qué un objeto Parquet sin extensión .parquet sigue siendo un Parquet válido, y cómo lo comprueban los lectores de esta sesión?
  • ¿Qué diferencia hay entre que un fichero fluya de HDFS a S3 como un flujo (de hdfs dfs -cat a upload_fileobj, como en esta sesión) y que se copie primero al disco del host?
  • ¿Qué parte de la configuración de los cuatro programas cambiaría, y qué parte no, al apuntar a un bucket real de Amazon S3 en vez de a RustFS?
  • ¿Por qué s3fs/PyArrow/Polars y httpfs/DuckDB pueden dar exactamente el mismo resultado numérico usando dos caminos de acceso distintos?
  • ¿Qué le falta a esta sesión para poder llamarse lakehouse en vez de data lake sobre S3?

Evidencias para la siguiente sesión¶

Antes de la siguiente sesión tendrás una reunión individual breve con el profesor para revisar el trabajo de esta sesión. Esa reunión combina una demostración en vivo sobre tu propio ordenador y una memoria escrita breve. Mantén RustFS y su bucket disponibles hasta entonces.

Qué mostrar en el ordenador durante la reunión¶

Ten preparado y a mano, funcionando en tu propio equipo:

  1. El bucket tcdm-datalake con los objetos bajo raw/tpcds/, en la consola de RustFS (http://localhost:9001) o en la salida de la celda de copia.
  2. El final de esa celda de copia: el número de piezas Parquet copiadas de las 24 tablas.
  3. La salida de check_s3_api.py, terminada en TCDM_S3_API_OK: la firma PAR1 leída con una petición por rango, sin descargar el objeto.
  4. Las tres líneas TCDM_SUMMARY (pyarrow-s3, polars-s3 y duckdb-s3) de date_dim, iguales entre sí y a las que obtuviste en S2 leyendo desde HDFS (73.049 filas).
  5. Que sabes señalar qué variables AWS_* habría que cambiar para usar un bucket real de Amazon S3 y cuáles no.

Memoria escrita (una o dos páginas)¶

Trae también un documento breve —una o dos páginas, no hace falta más— que no se limite a pegar capturas de las órdenes anteriores: debe explicar con tus propias palabras la diferencia entre un bloque HDFS y un objeto S3; qué son un bucket, una clave y un prefijo, y por qué un prefijo no es un directorio; por qué los tres lectores obtienen el mismo resultado que en S2 aunque haya cambiado el almacenamiento; y qué le falta a este data lake para poder llamarse lakehouse.

Qué es importante de cara al examen final¶

El examen no pide recordar la sintaxis exacta de una orden. Debes poder explicar:

  • qué ofrece un almacenamiento de objetos —objetos identificados por una clave dentro de un bucket, operaciones HTTP, lecturas por rango— frente a un sistema de ficheros distribuido como HDFS, con bloques, réplicas y un espacio de nombres;
  • por qué el formato Parquet es independiente de dónde se guarden los ficheros, y cómo se reconoce un Parquet sin fiarse de su extensión;
  • qué papel tienen boto3, s3fs y la extensión httpfs de DuckDB, y por qué caminos distintos devuelven el mismo dato;
  • qué hay que parametrizar —endpoint, credenciales, región y estilo de direccionamiento— para que un mismo programa funcione con RustFS y con Amazon S3.

Continuación del curso¶

Con HDFS y el almacenamiento de objetos sirviendo el mismo dataset TPC-DS SF1 en el mismo formato Parquet, la sesión 4 introduce Spark como motor de cómputo distribuido, leyendo los Parquet de HDFS. Las sesiones 5 a 8 añaden sobre HDFS el particionado físico, un catálogo (Hive Metastore y Trino), tablas Iceberg y su alimentación continua. Lo aprendido aquí sobre buckets, claves y prefijos es lo que permitiría llevar ese mismo recorrido a un almacenamiento de objetos.