Módulo 7: Catalogs Maintenance And Delta Lake By Contrast

Compactando archivos pequeños

Descripción

La lección 3 midió un problema de snapshots viejos: escrituras redundantes que ya nadie necesita, pero que siguen ocupando espacio hasta que alguien las expira. Esta lección mide un problema distinto, que puede existir incluso en una tabla sin ningún snapshot redundante: archivos pequeños dentro del snapshot vigente. Vas a cargar la misma semana de 40 órdenes de Kiosko de siempre, pero día por día en vez de en un solo append() — exactamente como llegaría en un pipeline real que recibe un archivo por día—, y vas a confirmar que eso, por sí solo, fragmenta la tabla en más archivos de los que hacen falta.

Conexión con el módulo. Esta lección no toca kiosko.dim_product — usa una tabla nueva, kiosko.fact_orders_daily_batches, siguiendo la misma disciplina de tablas dedicadas que ya viste en el módulo 6 (dim_product_upsert_demo). El motivo: el problema de esta lección es ortogonal al de la lección 3 — ninguna de las siete escrituras que vas a hacer acá es redundante, cada una trae datos de negocio reales y nuevos, y aun así el resultado es una tabla fragmentada.

Por qué esta lección usa fact_orders, no dim_product

kiosko.dim_product tiene cuatro filas — sin importar cuántas veces la reescribas, cada overwrite() cabe, siempre, en un solo archivo Parquet pequeño. El problema de "archivos pequeños" no aparece ahí de forma natural. kiosko.fact_orders, en cambio, es el tipo de tabla donde este problema aparece todo el tiempo en producción: recibe datos con cierta frecuencia —una vez al día, una vez por hora, a veces en streaming—, y cada llegada se convierte, típicamente, en su propio archivo. Esta lección reconstruye esa situación con los mismos datos de siempre, agrupados por el día real en que ocurrió cada orden.

Una analogía: el mismo álbum, ahora con una foto por visita en vez de un rollo completo

El estante del supermercado, en los módulos 1 y 3, siempre se reponía completo, así que cada foto capturaba el estante entero. Imagina, en cambio, un cliente que entra siete veces distintas durante la semana, y el empleado del control de calidad toma una foto nueva cada vez que ese cliente se retira, aunque solo haya movido dos o tres productos. Al final de la semana tienes siete fotos parciales, cada una válida y necesaria —ninguna es redundante, cada una documenta algo real que pasó—, pero juntarlas para responder "¿cómo se ve el estante completo hoy?" exige mirar las siete, una por una, en vez de una sola foto completa. Compactar, en este contexto, no es "borrar fotos" —eso sería perder información—, es reimprimir el mismo contenido en menos páginas: la misma información, organizada de forma más eficiente de leer.

Ejemplo trabajado: siete llegadas diarias, siete archivos vivos

Paso 1 — Carga la semana de Kiosko, un append() por día

# fact_orders_daily_batches.py -- la misma semana de 40 ordenes, cargada dia por dia
import os
from datetime import datetime

import pyarrow as pa
from pyiceberg.catalog import load_catalog
from pyiceberg.schema import Schema
from pyiceberg.types import DoubleType, IntegerType, NestedField, StringType, TimestampType

warehouse_path = os.path.abspath("kiosko_warehouse")
catalog_db_path = os.path.abspath("kiosko_catalog.db")

catalog = load_catalog(
    "kiosko", type="sql",
    uri=f"sqlite:///{catalog_db_path}", warehouse=f"file://{warehouse_path}",
)

# RAW_ORDERS_BY_DAY -- las mismas 40 filas del modulo 1, leccion 6, agrupadas por
# el dia real en que ocurrieron (2026-08-03 .. 2026-08-09, 7 dias, 8+6+2+5+7+9+3=40)
from raw_orders_by_day import RAW_ORDERS_BY_DAY

FACT_ORDERS_SCHEMA = Schema(
    NestedField(field_id=1, name="order_id", field_type=StringType(), required=True),
    NestedField(field_id=2, name="store_id", field_type=StringType(), required=True),
    NestedField(field_id=3, name="product_id", field_type=StringType(), required=True),
    NestedField(field_id=4, name="quantity", field_type=IntegerType(), required=True),
    NestedField(field_id=5, name="unit_price", field_type=DoubleType(), required=True),
    NestedField(field_id=6, name="revenue", field_type=DoubleType(), required=True),
    NestedField(field_id=7, name="order_ts", field_type=TimestampType(), required=True),
)
pa_schema = pa.schema([
    pa.field("order_id", pa.string(), nullable=False),
    pa.field("store_id", pa.string(), nullable=False),
    pa.field("product_id", pa.string(), nullable=False),
    pa.field("quantity", pa.int32(), nullable=False),
    pa.field("unit_price", pa.float64(), nullable=False),
    pa.field("revenue", pa.float64(), nullable=False),
    pa.field("order_ts", pa.timestamp("us"), nullable=False),
])

table = catalog.create_table("kiosko.fact_orders_daily_batches", schema=FACT_ORDERS_SCHEMA)

for day, orders in RAW_ORDERS_BY_DAY:
    rows = []
    for order_id, store_id, product_id, quantity, unit_price, ts in orders:
        rows.append({
            "order_id": order_id, "store_id": store_id, "product_id": product_id,
            "quantity": quantity, "unit_price": unit_price,
            "revenue": round(quantity * unit_price, 10),
            "order_ts": datetime.fromisoformat(ts),
        })
    table.append(pa.Table.from_pylist(rows, schema=pa_schema))
    print(f"append() del {day}: {len(rows)} filas")

Qué esperar (verificado corriendo el script real):

append() del 2026-08-03: 8 filas
append() del 2026-08-04: 6 filas
append() del 2026-08-05: 2 filas
append() del 2026-08-06: 5 filas
append() del 2026-08-07: 7 filas
append() del 2026-08-08: 9 filas
append() del 2026-08-09: 3 filas

Siete llamadas a append(), cada una con las órdenes reales de un solo día — nada redundante, cada fila es información de negocio genuina, exactamente como llegaría de un sistema transaccional real que exporta un archivo por día.

Paso 2 — Confirma el total, y cuenta los archivos vivos

total_rows = table.scan().to_arrow().num_rows
fact_rows = table.scan().to_arrow().to_pylist()
total_revenue = round(sum(r["revenue"] for r in fact_rows), 2)
print(f"\ntotal filas: {total_rows}, revenue total: {total_revenue}")

live_files = table.inspect.files()
print("table.inspect.files().num_rows:", live_files.num_rows)
for row in live_files.select(["file_path", "record_count", "file_size_in_bytes"]).to_pylist():
    print(" ", row["file_path"].split("/")[-1], "record_count=", row["record_count"],
          "bytes=", row["file_size_in_bytes"])

Qué esperar (verificado corriendo el script real; los nombres de archivo son de tu propia corrida, distintos cada vez — record_count y file_size_in_bytes son deterministas):

total filas: 40, revenue total: 106.15

table.inspect.files().num_rows: 7
  <archivo-1>.parquet record_count= 3 bytes= 2777
  <archivo-2>.parquet record_count= 9 bytes= 2928
  <archivo-3>.parquet record_count= 7 bytes= 2873
  <archivo-4>.parquet record_count= 5 bytes= 2834
  <archivo-5>.parquet record_count= 2 bytes= 2747
  <archivo-6>.parquet record_count= 6 bytes= 2857
  <archivo-7>.parquet record_count= 8 bytes= 2873

El revenue total —106.15— coincide, otra vez, con el número que ya verificaron las seis guías anteriores del ecosistema: cargar los datos día por día en vez de en un solo append() no cambió ni un centavo del resultado de negocio. Lo que sí cambió es la forma física de la tabla: siete archivos Parquet, todos vivos, todos necesarios para responder una consulta sobre la tabla completa — ninguno es candidato a expire_snapshots (la técnica de la lección 5 no aplica acá, porque no hay ningún snapshot redundante que expirar) — y cada uno diminuto: entre 2 y 9 filas, entre 2,7 y 2,9 KB.

Por qué esto es un problema, aunque ningún archivo sea redundante

La documentación oficial de Apache Iceberg lo resume así: "More data files leads to more metadata stored in manifest files, and small data files causes an unnecessary amount of metadata and less efficient queries from file open costs." Escanear kiosko.fact_orders_daily_batches completa —incluso un SELECT * trivial— obliga al motor de consulta a abrir siete archivos separados en vez de uno solo. Con siete archivos de unos pocos KB cada uno, el costo relativo de abrir cada archivo (ubicarlo en el almacenamiento, leer su footer Parquet, planear la lectura) puede superar, por mucho, el costo de leer los datos mismos — exactamente el motivo por el que un archivo de 500 MB casi siempre se lee más eficiente que cien archivos de 5 MB, aunque el volumen total de datos sea idéntico.

Compactar es la operación que resuelve esto: combinar varios archivos pequeños en menos archivos, más grandes, sin cambiar ni una fila de contenido — la estrategia por defecto de Iceberg para esto se llama bin-packing ("empaquetado de contenedores"): agrupa archivos pequeños hasta acercarse a un tamaño objetivo (target-file-size-bytes, 512 MB por defecto en la configuración de escritura de Iceberg), y escribe ese grupo como un solo archivo nuevo.

La operación real: rewrite_data_files — no disponible en PyIceberg 0.11.1 puro Python

Se verificó, contra el código fuente instalado de PyIceberg 0.11.1 y contra su referencia de API oficial, que no existe ningún método equivalente a la compactación en Table ni en table.maintenance — la única operación documentada bajo table.maintenance en esta versión es expire_snapshots() (lección 5). La compactación de archivos de datos, en el ecosistema Iceberg de 2026, es una operación que corre en paralelo sobre un motor de cómputo distribuido — típicamente Spark, a través de la acción rewriteDataFiles o el procedimiento SQL rewrite_data_files —, y ningún cliente Python puro la implementa todavía.

Este es exactamente el patrón que la lección 3 del módulo 6 ya documentó con MERGE INTO: la sintaxis existe, está verificada contra la documentación oficial, pero no corre en este entorno. El siguiente bloque está marcado como representativo:

-- (representativo) -- sintaxis verificada contra la documentacion oficial de Iceberg,
-- NO ejecutada en este entorno: PyIceberg 0.11.1 no implementa rewrite_data_files.

-- compactar kiosko.fact_orders_daily_batches con la estrategia bin-pack por defecto
CALL local.system.rewrite_data_files('kiosko.fact_orders_daily_batches');

-- version explicita, fijando un tamano objetivo mas chico (los defaults de Iceberg
-- estan pensados para tablas reales de produccion, no para 40 filas de laboratorio)
CALL local.system.rewrite_data_files(
    table => 'kiosko.fact_orders_daily_batches',
    options => map('target-file-size-bytes', '134217728')  -- 128 MB
);

Qué esperar (representativo): los siete archivos de entre 2,7 y 2,9 KB se combinarían en un solo archivo Parquet con las mismas 40 filas — el mismo revenue total (106.15), la misma estructura, ningún dato perdido ni duplicado—, y table.inspect.files().num_rows pasaría de 7 a 1. El procedimiento, corrido sobre Spark, produce además un reporte estructurado (rewritten_data_files_count, added_data_files_count, rewritten_bytes_count) que confirma exactamente cuántos archivos entraron y cuántos salieron — la misma disciplina de "verificar, no asumir" que esta guía aplicó en cada módulo anterior, ahora aplicada a una operación que no se ejecutó en este entorno.

Profundización: compactar produce un snapshot nuevo, no lo evita

Vale la pena una precisión que se malinterpreta seguido: compactar no reduce el número de snapshots — al contrario, agrega uno más. El tipo de operación que Iceberg registra para una compactación es replace (distinto de append, overwrite o delete): la especificación formal de Iceberg lo describe como "Data and delete files were added and removed without changing table data" — el contenido lógico de la tabla no cambia ni una fila, pero sí cambia qué archivos físicos representan ese contenido, y eso, como cualquier cambio, se archiva como un snapshot nuevo. Esta es la razón exacta por la que la lección 5 —expire_snapshots— y esta lección resuelven problemas relacionados pero distintos: compactar limpia el problema de "demasiados archivos pequeños en el snapshot vigente"; expirar limpia el problema de "demasiados snapshots viejos que ya nadie necesita" — y, sin querer, un rewrite_data_files sin ningún expire_snapshots posterior dejaría los siete archivos originales todavía rastreados por el snapshot anterior a la compactación, sin liberar ni un byte hasta que ese snapshot también se expire.

Diagrama: siete archivos vivos, un candidato a compactar

flowchart LR
    D1["2026-08-03\n8 filas"] --> F["kiosko.fact_orders_daily_batches\nsnapshot vigente"]
    D2["2026-08-04\n6 filas"] --> F
    D3["2026-08-05\n2 filas"] --> F
    D4["2026-08-06\n5 filas"] --> F
    D5["2026-08-07\n7 filas"] --> F
    D6["2026-08-08\n9 filas"] --> F
    D7["2026-08-09\n3 filas"] --> F

    F -->|"7 archivos vivos,\nninguno redundante"| Q["Costo: 7 aperturas de archivo\npor cada consulta completa"]
    Q -.->|"rewrite_data_files\n(representativo)"| C["1 archivo compactado\nmismas 40 filas\n+1 snapshot operation=replace"]

Errores comunes

Confundir el problema de esta lección con el de la lección 3. Qué pasa: alguien, después de ver "7 archivos" en esta lección y "7 archivos" en la lección 3 (una coincidencia numérica, no una relación causal), asume que expire_snapshots resolvería también el problema de los archivos pequeños de esta lección. Por qué pasa: ambos números son "7", y ambos son sobre archivos Parquet, así que es fácil mezclarlos. Cómo detectarlo: revisa si los archivos en cuestión son todos vivos (como en esta lección, ninguno redundante) o si algunos son redundantes (como en la lección 3, seis de siete sin valor de negocio adicional). Cómo corregirlo: expire_snapshots (lección 5) resuelve archivos que ya nadie necesita, dejados atrás por snapshots viejos; rewrite_data_files/compactación (esta lección) resuelve archivos que se necesitan, pero que están fragmentados en más piezas de las convenientes. Son ejes distintos, y una tabla real puede tener ambos problemas al mismo tiempo, como sería el caso si el pipeline nocturno de la lección 3 también hubiera llegado en batches diarios pequeños.

Asumir que un target-file-size-bytes más chico siempre es mejor porque "genera archivos más manejables". Qué pasa: alguien, viendo el ejemplo de esta lección con 128 MB, concluye que archivos más pequeños son, en general, más seguros o más fáciles de trabajar. Por qué pasa: en un dataset de laboratorio de 40 filas, cualquier tamaño objetivo por encima de unos pocos KB produce el mismo resultado —un solo archivo—, así que la intuición de "más chico es más manejable" no se pone a prueba. Cómo detectarlo: si tu tabla real tiene millones de filas y configuras un target-file-size-bytes demasiado pequeño, vas a terminar con muchos archivos "compactados" que siguen siendo, en términos relativos, pequeños — el mismo problema de file-open cost que esta lección describe, solo que después de haber gastado el trabajo de compactar. Cómo corregirlo: el valor por defecto de Iceberg (512 MB, la misma constante que aparece en la configuración de escritura de la tabla) es un punto de partida razonable para datasets de escala real — ajustarlo hacia abajo tiene sentido solo si tu patrón de consulta filtra agresivamente por partición y prefieres archivos más chicos por partición, un tema que profundiza cost-optimization-caching-guide, no esta lección.

Ejercicios

Ejercicio 1 — Reproduce el experimento completo tú mismo, y confirma los dos números. En un directorio nuevo, corre el script de esta lección. Confirma 40 filas totales, 106.15 de revenue, y 7 archivos vivos.

Ver solución

Si tu entorno tiene PyIceberg 0.11.1 instalado, tu salida debería coincidir en los tres números con esta lección. Los nombres de archivo (file_path) van a ser distintos —cada uno incluye un UUID generado en el momento de escribir—, pero record_count (3, 9, 7, 5, 2, 6, 8, en algún orden) y file_size_in_bytes deberían coincidir, porque dependen únicamente del contenido de cada archivo, no del momento en que corriste el script.

Ejercicio 2 — Calcula cuántos archivos habría producido esta misma semana si, en vez de agrupar por día, el pipeline hubiera agrupado por hora. Usando los timestamps de raw_orders_by_day (heredados del módulo 1, lección 6), cuenta cuántas horas distintas (order_ts truncado a la hora) tienen al menos una orden.

Ver solución

Contando los timestamps únicos por hora en las 40 órdenes de Kiosko, el resultado son bastantes más de siete horas distintas a lo largo de la semana —cada día tiene entre 2 y 4 horas distintas con al menos una orden—, así que agrupar por hora en vez de por día produciría más de veinte archivos, cada uno con un puñado de filas todavía menor que los de esta lección. Este ejercicio ilustra la relación directa entre frecuencia de escritura y fragmentación: cuanto más seguido escribes, sin ningún mecanismo de compactación corriendo detrás, más archivos pequeños acumulas — el motivo exacto por el que las tablas alimentadas por streaming (fuera del alcance de esta guía, terreno de streaming-with-kafka-and-flink-guide) casi siempre necesitan compactación programada como parte normal de su operación, no como una excepción.

Ejercicio 3 — Explica, en tus propias palabras, por qué rewrite_data_files produce una operación de tipo replace y no overwrite. Piensa en la diferencia entre "cambiar el contenido lógico de una tabla" y "cambiar cómo ese mismo contenido está organizado físicamente".

Ver solución

overwrite() (módulo 3) cambia el contenido lógico: las filas que la tabla devuelve antes y después de la operación son distintas —P002 pasa de snacks a health-snacks, por ejemplo—. replace, el tipo de operación que produce una compactación, no cambia ni una fila de lo que la tabla devuelve —las mismas 40 órdenes, con los mismos valores, siguen ahí—; lo único que cambia es cuántos archivos físicos representan ese mismo contenido. Marcar esta operación con un tipo distinto (replace, no overwrite) le permite a cualquier herramienta que lea el historial de la tabla —incluida la propia lógica interna de expire_snapshots— distinguir, sin ambigüedad, entre "acá cambió algo que un usuario de negocio necesita saber" y "acá solo se reorganizaron los archivos, el contenido es idéntico al del snapshot anterior".

Resumen y siguiente paso

En esta lección identificaste un problema de mantenimiento distinto al de la lección 3: no snapshots redundantes, sino archivos pequeños dentro de un snapshot vigente y completamente legítimo. Cargaste kiosko.fact_orders_daily_batches día por día —siete append() reales, ninguno redundante— y confirmaste, con table.inspect.files(), que el resultado son siete archivos diminutos, entre 2,7 y 2,9 KB cada uno. Verificaste que PyIceberg 0.11.1 no implementa rewrite_data_files en Python puro, y documentaste la sintaxis representativa de Spark que sí lo resuelve, incluida la precisión de que compactar produce un snapshot nuevo (operation=replace), no lo evita.

Antes de avanzar deberías poder: distinguir el problema de esta lección del de la lección 3; y explicar por qué compactar sin expirar después deja el espacio de los archivos viejos sin liberar.

La lección 5 vuelve a kiosko.dim_product y a los trece snapshots de la lección 3, y esta vez sí ejecuta de verdad: table.maintenance.expire_snapshots(), con snap_v1 explícitamente protegido.

Recursos

  • Apache Iceberg — documentación oficial, "Maintenance", sección "Compact data files", fuente de la cita textual sobre el costo de los archivos pequeños. iceberg.apache.org/docs/latest/maintenance. En inglés.
  • Apache Iceberg — documentación oficial, "Spark Procedures", sección rewrite_data_files, fuente exacta de la sintaxis CALL representativa de esta lección. iceberg.apache.org/docs/latest/spark-procedures. En inglés.
  • PyIceberg — referencia de API, la sección table.maintenance, confirmando que solo expire_snapshots está documentado en esta versión. py.iceberg.apache.org/api. En inglés.
  • DISEÑO de esta guía — la sección del módulo 7, "compactación de archivos pequeños". src/guides/lakehouse-and-iceberg-guide/DISENO.md. En español.