Estoy tratando de dividir un conjunto de archivos S3 como se muestra a continuación en función de una columna en carpetas basadas en columnas individuales. No estoy seguro del problema con mi código a continuación.
column 1, column 2 20130401, value1 20130402, value2 20130403, value3 val newDataDF = sqlContext.read.parquet("s3://xxxxxxx-bucket/basefolder/") newDataDF.cache() val uniq_days = newDataDF.select(newDataDF("column1")).distinct.show() uniq_days.cache() uniq_days.foreach(x => {newDataDF.filter(newDataDF("column1") === x).write.save(s"s3://xxxxxx-bucket/partitionedfolder/$x/")})¿Puedes ayudarme? Incluso una versión pyspark está bien. Estoy buscando la siguiente estructura.
s3://xxxxxx-bucket/partitionedfolder/20130401/part-*** column 1, column 2 20130401, value 1 s3://xxxxxx-bucket/partitionedfolder/20130402/part-*** column 1, column 2 20130402, value 1 s3://xxxxxx-bucket/partitionedfolder/20130403/part-*** column 1, column 2 20130403, value 1aquí está el error
org.apache.spark.SparkException: Job aborted due to stage failure: Task 22 in stage 82.0 failed 4 times, most recent failure: Lost task 22.3 in stage 82.0 (TID 2753 Driver stacktrace: at org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1454) Caused by: java.lang.NullPointerExceptionActualizar con la solución actual:
val newDataDF = sqlContext.read.parquet("s3://xxxxxx-bucket/basefolder/") newDataDF.cache() val uniq_days = newDataDF.select(newDataDF("column1")).distinct.rdd.map(_.getString(0)).collect().toList uniq_days.foreach(x => {newDataDF.filter(newDataDF("column1") === x).write.save(s"s3://xxxxxx-bucket/partitionedfolder/$x/")})Creo que te perdiste "s" en el guardado. :)
http://docs.scala-lang.org/overviews/core/string-interpolation.html#the-s-string-interpolator
Cambio:
write.save("s3://xxxxxx-bucket/partitionedfolder/$x/")})A:
write.save(s"s3://xxxxxx-bucket/partitionedfolder/$x/")})Hay más problemas, show nunca devuelve ningún valor.
Cambio:
val uniq_days = newDataDF.select(newDataDF("mevent_day")).distinct.show() uniq_days.cache()A:
val uniq_days = newDataDF.select(newDataDF("mevent_day")).distinct.rdd.map(_.getString(0)).collect().toList