Empresas
Empleos
  • Sobre nosotros
  • Soluciones
    • Publicación de vacantes
      Publica tu vacante y recibe candidatos calificados en 48h.
    • Evaluación de candidatos
      500+ pruebas técnicas y psicológicas, más anti-fraude.
    • Headhunting
      Búsqueda ejecutiva a la medida de principio a fin.
    • Nómina + EOR
      Dispersión de nómina y EOR en más de 15 países de LATAM.
  • Precios
  • Empleos

0

124
Vistas
Spark divide un archivo en varias carpetas en función de un campo

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 1

aquí 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.NullPointerException

Actualizar 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/")})
about 4 years ago · Santiago Trujillo
1 Respuestas
Responde la pregunta

0

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
about 4 years ago · Santiago Trujillo Denunciar
Responde la pregunta
Encuentra empleos remotos

¡Descubre la nueva forma de encontrar empleo!

Top de empleos
Top categorías de empleo
Empresas
Publicar vacante Precios Comercial
Legal
Términos y condiciones Política de privacidad
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomiéndame algunas ofertas
Necesito ayuda