Spark
JuniorEl 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.
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.
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.