Spark

Junior

El problema de los archivos pequeños: por qué tu job de Spark pasa más tiempo en I/O que en cómputo

Un job puede estar limitado por I/O y parecer limitado por cómputo en cualquier dashboard que solo mida CPU, porque el cuello de botella está en abrir archivos, no en leer bytes.

Leer en:
Overhead de apertura Overhead de planificación Tiempo de lectura real
A dónde va el tiempo en un stage afectado por archivos pequeños, aproximadamente - ilustrativo, no medido.

Por debajo de cierto tamaño de archivo, el costo fijo de abrirlo - búsqueda de metadatos, planificación de la tarea, configuración del códec - cuesta más que leer su contenido real.

Cómo terminas con miles de archivos diminutos

Las fuentes habituales: particionar una escritura por una columna de alta cardinalidad (cada valor de partición obtiene su propio archivo pequeño), jobs de streaming que confirman micro-batches pequeños con un intervalo de disparo corto, o un stage de shuffle con muchas más particiones de salida de las que los datos realmente necesitan. Ninguna de estas es exactamente un error - a menudo es la decisión correcta para el job que escribe los datos - pero deja un desorden para quien lea a continuación.

Por qué cuesta más de lo que sugiere el conteo de bytes

Cada archivo, por pequeño que sea, sigue necesitando una tarea completa: planificarla, abrir el archivo (un round trip de red real en S3 o GCS), leer un footer o header, y luego cerrarlo. Con 10.000 archivos diminutos, ese overhead por archivo puede dominar el tiempo total del stage aunque los datos reales apenas llenen la memoria de unos pocos executors - el job parece limitado por cómputo en un gráfico de CPU porque el CPU está ocupado planificando y gestionando I/O, no haciendo trabajo real.

# tamaño medio aproximado de archivo para una ruta de salida dada avg_file_size = total_bytes_written / file_count # notablemente por debajo de ~128MB (Parquet) vale la pena investigarlo

El arreglo que realmente funciona

coalesce() o repartition() justo antes de la escritura, apuntando a aproximadamente 128MB-1GB por archivo de salida según el formato y el motor que lo lea después. Un arreglo único en el lado de la escritura es más barato que cada job posterior pagando repetidamente el impuesto de abrir archivos - compactar después también funciona, pero es un costo recurrente que una escritura correctamente dimensionada evita por completo.

Mira cómo se ve esto en tu propio pipeline.

Sube un event log real de Spark, un run_results.json de dbt, o una exportación de métricas de Flink, y recibe recomendaciones concretas que requieren tu aprobación - no otra regla general.