【问题标题】:extract value from a list of json in pyspark从 pyspark 中的 json 列表中提取值
【发布时间】:2021-09-22 23:45:06
【问题描述】:

我有一个数据框,其中一列采用 json 列表的形式。我想从列中提取特定值(分数)并创建独立列。

raw_data = [{"user_id" : 1234, "col" : [{"id":14577120145280,"score":64.71,"Elastic_position":0},{"id":14568530280240,"score":88.53,"Elastic_position":1},{"id":14568530119661,"score":63.75,"Elastic_position":2},{"id":14568530205858,"score":62.79,"Elastic_position":3},{"id":14568530414899,"score":60.88,"Elastic_position":4}]}]

df = pd.DataFrame.from_dict(raw_data)

我想将我的结果数据框分解为:

【问题讨论】:

    标签: python pandas list pyspark


    【解决方案1】:

    假设你的 json 看起来像这样

    # a.json
    # {
    #     "user_id" : 1234,
    #     "col" : [
    #         {"id":14577120145280,"score":64.71,"Elastic_position":0},
    #         {"id":14568530280240,"score":88.53,"Elastic_position":1},
    #         {"id":14568530119661,"score":63.75,"Elastic_position":2},
    #         {"id":14568530205858,"score":62.79,"Elastic_position":3},
    #         {"id":14568530414899,"score":60.88,"Elastic_position":4}
    #     ]
    # }
    

    您可以阅读它,将其展平,然后像这样旋转它

    from pyspark.sql import functions as F
    from pyspark.sql import types as T
    
    schema = T.StructType([
        T.StructField('user_id', T.IntegerType()),
        T.StructField('col', T.StringType()),
    ])
    
    df = spark.read.json('a.json', multiLine=True, schema=schema)
    df.show(10, False)
    
    # +-------+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
    # |user_id|col                                                                                                                                                                                                                                                                                           |
    # +-------+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
    # |1234   |[{"id":14577120145280,"score":64.71,"Elastic_position":0},{"id":14568530280240,"score":88.53,"Elastic_position":1},{"id":14568530119661,"score":63.75,"Elastic_position":2},{"id":14568530205858,"score":62.79,"Elastic_position":3},{"id":14568530414899,"score":60.88,"Elastic_position":4}]|
    # +-------+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
    
    
    df.printSchema()
    # root
    #  |-- user_id: integer (nullable = true)
    #  |-- col: string (nullable = true)
    
    (df
        # this will parse your JSON string to JSON object
        .withColumn('col', F.from_json(
            F.col('col'),
            T.ArrayType(T.StructType([
                T.StructField('id', T.LongType()),
                T.StructField('score', T.DoubleType()),
                T.StructField('Elastic_position', T.IntegerType()),
            ]))
        ))
     
        .select('user_id', F.explode('col'))
        .groupBy('user_id')
        .pivot('col.Elastic_position')
        .agg(F.first('col.score'))
        .show(10, False)
    )
    
    # Output
    # +-------+-----+-----+-----+-----+-----+
    # |user_id|0    |1    |2    |3    |4    |
    # +-------+-----+-----+-----+-----+-----+
    # |1234   |64.71|88.53|63.75|62.79|60.88|
    # +-------+-----+-----+-----+-----+-----+
    

    【讨论】:

    • 列“col”的数据类型是字符串(而不是json)。所以我无法爆炸列。有什么建议吗?
    • 我更新了从字符串解析JSON的答案,你可以试试吗?
    【解决方案2】:

    尝试将pd.Series.explodegroupby 一起使用:

    df = pd.DataFrame.from_dict(raw_data).explode('col')
    df.assign(col=df['col'].str['score']).groupby('user_id').agg(list).apply(lambda x: (y:=x.explode()).set_axis(y.index + '_' + y.groupby(level=0).cumcount().astype(str)), axis=1).reset_index()
    

       user_id  col_0  col_1  col_2  col_3  col_4
    0     1234  64.71  88.53  63.75  62.79  60.88
    

    如果首先构造一个数据框并分解col 列,然后按重复的user_ids 分组并执行另一个explode 以使其长到宽,然后将前缀0 添加到4cumcount.

    【讨论】:

    • 非常感谢!您能否建议一种使用 pyspark 的方法?我的数据框太大,所以在 pandas 中处理是一个挑战。
    • @Teresa 哦,好吧,我从来没有在 pyspark 中编写过代码,但至少点赞会有帮助 :) 也许你应该尝试在 pyspark 中实现相同的逻辑。
    • 你能复习一下语法吗?我遇到了错误。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-09-16
    • 1970-01-01
    • 2019-07-29
    • 1970-01-01
    • 2019-10-20
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多