cómo agregar varias columnas y generar una cadena formateada en Python
Oct 29 2020
Tengo una tabla / marco de fecha con la siguiente estructura
**userid item createdon**
user1 item1 2020-10-01
user1 item2 2020-10-02
user1 item3 2020-10-03
user2 item4 2020-01-01
user2 item1 2020-03-03
...
para cada ID de usuario, necesito generar una cadena en un formato como este:
para user1:
{ "date": "2020-10-01",
"item": {"display": "item1"}
},
{ "date": "2020-10-02",
"item": {"display": "item2"}
},
{ "date": "2020-10-03",
"item": {"display": "item3"}
}
para el usuario2:
{ "date": "2020-01-01",
"item": {"display": "item4"}
},
{ "date": "2020-03-03",
"item": {"display": "item1"}
}
Estoy usando pyspark. Me pregunto si podría lograr esto utilizando la transformación de mapas. Cualquier ayuda sería apreciada.
Respuestas
1 user238607 Oct 30 2020 at 15:22
Podría hacer que el ejemplo funcione. Escribe un archivo por usuario.
Ref: Spark: escriba JSON varios archivos de DataFrame en función de la separación por valor de columna
from pyspark.sql import Row
from pyspark.sql.types import *
from pyspark import SparkContext, SQLContext
from pyspark.sql.functions import *
sc = SparkContext('local')
sqlContext = SQLContext(sc)
data1 = [
("user1", "item1", "2020-10-01"),
("user1", "item2", "2020-10-02"),
("user1", "item3", "2020-10-03"),
("user2", "item4", "2020-01-01"),
("user2", "item1", "2020-03-03")
]
df1Columns = ["userid", "item", "createdon"]
df1 = sqlContext.createDataFrame(data=data1, schema = df1Columns)
df1.printSchema()
df1.show(truncate=False)
partialSchema = StructType([StructField("date", StringType(), nullable=True)
, StructField("item", StructType([StructField("display", StringType(), nullable=True)]), nullable=True)
])
actualSchema = StructType([StructField("userid", StringType(), nullable=True)
, StructField("dict", partialSchema, nullable=True)
])
res_df = df1.rdd.map(lambda row: Row(row[0], {"date": row[2], "item" : {'display':row[1]}}))\
.toDF(actualSchema)
res_df.show(20, False)
res_df.repartition(col("userid")).select(col("userid"), col("dict.*")).write.partitionBy("userid").json("./helloworld/data/")
La última línea escribe dos archivos, uno por cada usuario.
Contenido del primer archivo de usuario:
{"date":"2020-10-01","item":{"display":"item1"}}
{"date":"2020-10-02","item":{"display":"item2"}}
{"date":"2020-10-03","item":{"display":"item3"}}
Contenido del segundo archivo de usuario:
{"date":"2020-01-01","item":{"display":"item4"}}
{"date":"2020-03-03","item":{"display":"item1"}}