【发布时间】:2021-01-12 21:35:05
【问题描述】:
我是 pyspark 的新手,我想以一种将每个值分配给新列的方式分解数组值。我尝试使用explode,但无法获得所需的输出。以下是我的输出
这是代码
from pyspark.sql import *
from pyspark.sql.functions import explode
if __name__ == "__main__":
spark = SparkSession.builder \
.master("local[3]") \
.appName("DataOps") \
.getOrCreate()
dataFrameJSON = spark.read \
.option("multiLine", True) \
.option("mode", "PERMISSIVE") \
.json("data.json")
dataFrameJSON.printSchema()
sub_DF = dataFrameJSON.select(explode("values.line").alias("new_values"))
sub_DF.printSchema()
sub_DF2 = sub_DF.select("new_values.*")
sub_DF2.printSchema()
sub_DF.show(truncate=False)
new_DF = sub_DF2.select("id", "period.*", "property")
new_DF.show(truncate=False)
new_DF.printSchema()
这是数据:
{
"values" : {
"line" : [
{
"id" : 1,
"period" : {
"start_ts" : "2020-01-01T00:00:00",
"end_ts" : "2020-01-01T00:15:00"
},
"property" : [
{
"name" : "PID",
"val" : "P120E12345678"
},
{
"name" : "EngID",
"val" : "PANELID00000000"
},
{
"name" : "TownIstat",
"val" : "12058091"
},
{
"name" : "ActiveEng",
"val" : "5678.1"
}
]
}
}
【问题讨论】:
-
请将代码示例、错误输出放在文本中而不是图像中,以便社区更容易为您提供帮助。
-
@Mikayel Saghyan,我希望您现在可以看到代码和示例数据?我正在尝试生成上面链接中给出的所需输出
标签: apache-spark pyspark apache-spark-sql pyspark-dataframes