como agregar várias colunas e gerar uma string formatada em python

Oct 29 2020

Eu tenho uma tabela / data com a seguinte estrutura

**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 usuário, preciso gerar uma string em um formato como este:

para o usuário1:

{ "date": "2020-10-01",
  "item": {"display": "item1"}
},
{ "date": "2020-10-02",
  "item": {"display": "item2"}
},
{ "date": "2020-10-03",
  "item": {"display": "item3"}
}

para o usuário2:

{ "date": "2020-01-01",
  "item": {"display": "item4"}
},
{ "date": "2020-03-03",
  "item": {"display": "item1"}
}

Estou usando o pyspark. Eu me pergunto se eu poderia conseguir isso utilizando a transformação do mapa? Qualquer ajuda seria apreciada.

Respostas

1 user238607 Oct 30 2020 at 15:22

Eu poderia fazer o exemplo funcionar. Ele grava um arquivo por usuário.

Ref: Spark: gravar vários arquivos JSON do DataFrame com base na separação por valor de coluna

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/")

A última linha grava dois arquivos, um para cada usuário.

Conteúdo do primeiro arquivo do usuário:

{"date":"2020-10-01","item":{"display":"item1"}}
{"date":"2020-10-02","item":{"display":"item2"}}
{"date":"2020-10-03","item":{"display":"item3"}}

Conteúdo do segundo arquivo do usuário:

{"date":"2020-01-01","item":{"display":"item4"}}
{"date":"2020-03-03","item":{"display":"item1"}}