【问题标题】:Explode array values into multiple columns using PySpark使用 PySpark 将数组值分解为多列
【发布时间】: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


【解决方案1】:

这是一个通用的解决方案,即使在 JSON 混乱的情况下也可以使用(元素的不同顺序或某些元素丢失)

您必须首先展平,regexp_replace 拆分“属性”列,最后是pivot。这也避免了新列名的硬编码。

构建您的数据框:

from pyspark.sql.types import *
from pyspark.sql import functions as F
from pyspark.sql.functions import col
from pyspark.sql.functions import *

schema = StructType([StructField("id", IntegerType()), StructField("start_ts", StringType()), StructField("end_ts", StringType()), \
    StructField("property", ArrayType(StructType(  [StructField("name", StringType()),  StructField("val", StringType())]    )))])

data = [[1, "2010", "2020", [["PID", "P123"], ["Eng", "PA111"], ["Town", "999"], ["Act", "123.1"]]],\
         [2, "2011", "2012", [["PID", "P456"], ["Eng", "PA222"], ["Town", "777"], ["Act", "234.1"]]]]

df = spark.createDataFrame(data,schema=schema)

df.show(truncate=False)
+---+--------+------+------------------------------------------------------+    
|id |start_ts|end_ts|property                                              |
+---+--------+------+------------------------------------------------------+
|1  |2010    |2020  |[[PID, P123], [Eng, PA111], [Town, 999], [Act, 123.1]]|
|2  |2011    |2012  |[[PID, P456], [Eng, PA222], [Town, 777], [Act, 234.1]]|
+---+--------+------+------------------------------------------------------+

展平和旋转:

df_flatten = df.rdd.flatMap(lambda x: [(x[0],x[1], x[2], y) for y in x[3]]).toDF(['id', 'start_ts', 'end_ts', 'property'])\
            .select('id', 'start_ts', 'end_ts', col("property").cast("string"))

df_split = df_flatten.select('id', 'start_ts', 'end_ts', regexp_replace(df_flatten.property, "[\[\]]", "").alias("replacced_col"))\
                .withColumn("arr", split(col("replacced_col"), ", "))\
                .select(col("arr")[0].alias("col1"), col("arr")[1].alias("col2"), 'id', 'start_ts', 'end_ts')

final_df = df_split.groupby(df_split.id,)\
                        .pivot("col1")\
                        .agg(first("col2"))\
                        .join(df,'id').drop("property")

输出:

final_df.show()
+---+-----+-----+----+----+--------+------+
| id|  Act|  Eng| PID|Town|start_ts|end_ts|
+---+-----+-----+----+----+--------+------+
|  1|123.1|PA111|P123| 999|    2010|  2020|
|  2|234.1|PA222|P456| 777|    2011|  2012|
+---+-----+-----+----+----+--------+------+

【讨论】:

  • 非常感谢您的回复,通过使用您的代码,我收到以下错误:TypeError: col() missing 1 required positional argument: 'strg',我认为 col 包有问题?你能分享你完整的代码,包括所有的进口吗?谢谢
  • 我已经添加了导入。你在哪里得到错误?
  • 非常感谢您的回复,它现在对我也有效
  • 为什么要麻烦而不是简单地通过访问相关数据来提取相关数据,然后创建新列?我不明白,这很复杂。
  • 这是因为这是一个通用的解决方案。我知道只对这个问题使用索引更简单。但是,如果 'property' 中元素的 order 发生变化,或者其中一个元素missing,则 indices 方法将不起作用。我在我的解决方案中考虑了这些 messy JSONs 的情况。 (并且没有对列名进行硬编码。适用于“n”个元素)
【解决方案2】:

您能否包含数据而不是屏幕截图?

同时,假设df是正在使用的数据框,我们需要做的是创建一个新的数据框,同时将vals从之前的property数组中提取到新列中,并删除@ 987654325@最后一栏:

from pyspark.sql.functions import col
output_df = df.withColumn("PID", col("property")[0].val).withColumn("EngID", col("property")[1].val).withColumn("TownIstat", col("property")[2].val).withColumn("ActiveEng", col("property")[3].val).drop("property")

如果 element 是 ArrayType 类型,请使用以下内容:

from pyspark.sql.functions import col
output_df = df.withColumn("PID", col("property")[0][1]).withColumn("EngID", col("property")[1][1]).withColumn("TownIstat", col("property")[2][1]).withColumn("ActiveEng", col("property")[3][1]).drop("property")

Explode 会将数组分解为新的行,而不是列,请参见:pyspark explode

【讨论】:

  • 感谢您的回复,我尝试了您建议的代码,但出现以下错误:TypeError: col() missing 1 required positional argument: 'strg'
  • 你能在上面的问题中包含你的代码吗?有数据样本?当你编辑时,试着编辑你自己的问题,而不是我的答案,你在正确的道路上
  • 代码和示例数据现已可用,请帮助我生成屏幕截图中提供的所需输出,提前谢谢
  • 您可能忘记在 col() 中插入“property”字符串,您可以在使用我的答案后显示所有代码吗?我刚刚在 Databricks 上使用了相同的代码,它工作得很好,没有错误
  • sub_DF = dataFrameJSON.select("UrbanDataset.values.line") sub_DF2 = dataFrameJSON.select(explode("UrbanDataset.values.line").alias("new_values")) sub_DF3 = sub_DF2。 select("new_values.*") new_DF = sub_DF3.select("id", "period.*", "property") new_DF.show(truncate=False) output_df = new_DF.withColumn("PID", col("property ")[0][1]) \ .withColumn("EngID", col("property")[1][1]) \ .withColumn("TownIstat", col("property")[2][1] ) \ .withColumn("ActiveEng", col("property")[3][1]).drop("property") output_df.show(truncate=False)
猜你喜欢
  • 1970-01-01
  • 2018-03-02
  • 2017-04-22
  • 2018-12-06
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多