Escribo para ver si alguien sabe cómo acelerar los tiempos de escritura de S3 desde Spark ejecutándose en EMR.
My Spark Job tarda más de 4 horas en completarse, sin embargo, el clúster solo está bajo carga durante las primeras 1,5 horas.
Tenía curiosidad por saber qué estaba haciendo Spark todo este tiempo. Miré los registros y encontré muchos comandos s3 mv , uno para cada archivo. Luego, al mirar directamente a S3, veo que todos mis archivos están en un directorio _temporal .
Secundario, me preocupa el costo de mi clúster, parece que necesito comprar 2 horas de cómputo para esta tarea específica. Sin embargo, termino comprando hasta 5 horas. Tengo curiosidad por saber si EMR AutoScaling puede ayudar con el costo en esta situación.
Algunos artículos discuten cómo cambiar el algoritmo de confirmación de salida del archivo, pero he tenido poco éxito con eso.
sc.hadoopConfiguration.set("mapreduce.fileoutputcommitter.algorithm.version", "2")Escribir en el HDFS local es rápido. Tengo curiosidad si emitir un comando hadoop para copiar los datos a S3 sería más rápido.
Lo que está viendo es un problema con outputcommitter y s3. el trabajo de confirmación aplica fs.rename en la carpeta _temporal y dado que S3 no admite el cambio de nombre, significa que una sola solicitud ahora está copiando y eliminando todos los archivos de _temporal a su destino final.
sc.hadoopConfiguration.set("mapreduce.fileoutputcommitter.algorithm.version", "2") solo funciona con la versión de hadoop > 2.7. lo que hace es copiar cada archivo de _temporal en la tarea de confirmación y no confirmar la tarea para que se distribuya y funcione bastante rápido.
Si usa una versión anterior de Hadoop, usaría Spark 1.6 y usaría:
sc.hadoopConfiguration.set("spark.sql.parquet.output.committer.class","org.apache.spark.sql.parquet.DirectParquetOutputCommitter")* tenga en cuenta que no funciona con la especulación activada o escribiendo en modo de adición
**También tenga en cuenta que está en desuso en Spark 2.0 (reemplazado por algoritmo.versión=2)
Por cierto, en mi equipo escribimos con Spark a HDFS y usamos trabajos DISTCP (específicamente s3-dist-cp) en producción para copiar los archivos a S3, pero esto se hace por varias otras razones (coherencia, tolerancia a fallas), por lo que no es necesario .. puedes escribir en S3 bastante rápido usando lo que sugerí.
Tuve un caso de uso similar en el que usé chispa para escribir en s3 y tuve un problema de rendimiento. La razón principal fue que Spark estaba creando muchos archivos parciales de cero bytes y reemplazar los archivos temporales por el nombre real del archivo estaba ralentizando el proceso de escritura. Intenté el siguiente enfoque como solución alternativa
Escriba la salida de Spark en HDFS y use Hive para escribir en s3. El rendimiento fue mucho mejor ya que Hive creaba menos archivos de piezas. El problema que tuve es (también tuve el mismo problema al usar Spark), la acción de eliminación en la política no se proporcionó en prod env por razones de seguridad. En mi caso, el cubo S3 estaba encriptado en kms.
Escriba la salida de chispa en HDFS y copie los archivos hdfs en local y use aws s3 copy para enviar datos a s3. Tuvo los segundos mejores resultados con este enfoque. Boleto creado con Amazon y sugirieron ir con este.
Use s3 dist cp para copiar archivos de HDFS a S3. Esto estaba funcionando sin problemas, pero sin rendimiento
El comunicador directo fue sacado de la chispa ya que no era resistente a las fallas. Recomiendo encarecidamente no usarlo.
Se está trabajando en Hadoop, s3guard, para agregar 0-rename committers, que serán O(1) y tolerantes a fallas; mantener un ojo en HADOOP-13786 .
Ignorando "el confirmador de Magic" por ahora, el confirmador de pruebas basado en Netflix se enviará primero (¿hadoop 2.9? ¿3.0?)
Resultado: la confirmación de la tarea toma segundos de datos/ancho de banda, pero la confirmación del trabajo no toma más tiempo que hacer 1-4 GET en la carpeta de destino y un POST para cada archivo pendiente, este último en paralelo.
Puede elegir el committer en el que se basa este trabajo, de netflix , y probablemente usarlo en Spark hoy. Establezca el algoritmo de confirmación del archivo = 1 (debería ser el valor predeterminado) o en realidad no escribirá los datos.