Tengo las siguientes rutas de archivo que leemos con particiones en s3
prefix/company=abcd/service=xyz/date=2021-01-01/file_01.json prefix/company=abcd/service=xyz/date=2021-01-01/file_02.json prefix/company=abcd/service=xyz/date=2021-01-01/file_03.jsonCuando leo esto con pyspark
self.spark \ .read \ .option("basePath", 'prefix') \ .schema(self.schema) \ .json(['company=abcd/service=xyz/date=2021-01-01/'])Todos los archivos tienen el mismo esquema y se cargan en la tabla como filas. Un archivo podría ser algo como esto:
{"id": "foo", "color": "blue", "date": "2021-12-12"} El problema es que a veces los archivos tienen el campo de fecha que choca con mi código de partición, como date . Entonces, quiero saber si es posible cargar los archivos sin las columnas de partición, cambiar el nombre de la columna de fecha JSON y luego agregar las columnas de partición.
La mesa final sería:
| id | color | file_date | company | service | date | ------------------------------------------------------------- | foo | blue | 2021-12-12 | abcd | xyz | 2021-01-01 | | bar | red | 2021-10-10 | abcd | xyz | 2021-01-01 | | baz | green | 2021-08-08 | abcd | xyz | 2021-01-01 |EDITAR:
Más información: tengo 5 o 6 particiones a veces y la fecha es una de ellas (no la última). Y también necesito leer varias particiones de fecha a la vez. El esquema que paso a Spark también contiene las columnas de partición, lo que lo hace más complicado.
No controlo los datos de entrada, así que necesito leerlos tal cual. Puedo cambiar el nombre de las columnas del archivo pero no las columnas de la partición.
¿Sería posible agregar un alias a las columnas del archivo como lo haríamos al unir 2 marcos de datos?
Chispa 3.1
Una forma es listar los archivos bajo la ruta del prefix S3 usando, por ejemplo, Hadoop FS API , luego pasar esa lista a spark.read . De esta manera, Spark no los detectará como particiones y podrá cambiar el nombre de las columnas del archivo si es necesario.
Después de cargar los archivos en un marco de datos, recorra las columnas df y cambie el nombre de las que también están presentes en su partitions_colums de columnas_particiones (agregando el prefijo del file por ejemplo).
Finalmente, analice las particiones desde input_file_name() usando la función regexp_extract .
Aquí hay un ejemplo:
from pyspark.sql import functions as F Path = sc._gateway.jvm.org.apache.hadoop.fs.Path conf = sc._jsc.hadoopConfiguration() s3_path = "s3://bucket/prefix" file_cols = ["id", "color", "date"] partitions_cols = ["company", "service", "date"] # listing all files for input path json_files = [] files = Path(s3_path).getFileSystem(conf).listFiles(Path(s3_path), True) while files.hasNext(): path = files.next().getPath() if path.getName().endswith(".json"): json_files.append(path.toString()) df = spark.read.json(json_files) # you can pass here the schema of the files without the partition columns # renaming file column in if exists in partitions df = df.select(*[ F.col(c).alias(c) if c not in partitions_cols else F.col(c).alias(f"file_{c}") for c in df.columns ]) # parse partitions from filenames for p in partitions_cols: df = df.withColumn(p, F.regexp_extract(F.input_file_name(), f"/{p}=([^/]+)/", 1)) df.show() #+-----+----------+---+-------+-------+----------+ #|color| file_date| id|company|service| date| #+-----+----------+---+-------+-------+----------+ #|green|2021-08-08|baz| abcd| xyz|2021-01-01| #| blue|2021-12-12|foo| abcd| xyz|2021-01-01| #| red|2021-10-10|bar| abcd| xyz|2021-01-01| #+-----+----------+---+-------+-------+----------+Lo más fácil sería simplemente cambiar el nombre de la columna de partición. A continuación, puede leer los datos y cambiar el nombre de las columnas como desee. Tampoco perderá los beneficios de la partición.
Si esa no es una opción, puede leer en los jsons usando un comodín para las particiones, cambie el nombre de la columna de fecha a 'file_date' y luego agregue la fecha de la partición extrayéndola del nombre del archivo. Puede obtener el nombre de archivo de input_file_name en pyspark.sql.functions .
Editar: me perdí de que tenga otras columnas particionadas antes de la fecha, también tendría que extraerlas del nombre del archivo para que sea menos ideal.
Sí, podemos leer todos los archivos json sin columnas de partición. Utilice directamente la ruta de la carpeta principal y cargará todos los datos de las particiones en el marco de datos.
Después de leer el marco de datos, puede usar la función withColumn() para cambiar el nombre del campo de fecha.
Algo como lo siguiente debería funcionar
df= spark.read.json("s3://bucket/table/**/*.json") renamedDF= df.withColumnRenamed("old column name","new column name")