【问题标题】:Converting a pyspark dataframe to a nested json object将 pyspark 数据框转换为嵌套的 json 对象
【发布时间】:2021-06-12 16:54:37
【问题描述】:

我有一个spark数据框如下

----------------------------------------------------------------------------
| item_id |   popular_tags   | popularity_score
____________________________________________________________________________
| id_1        Samsung         0.4
| id_1        long battery    0.8
| id_2        Apple           0.9
| id_2        UI              0.9
_____________________________________________________________________________

我想将此数据框按item_id 分组并输出为一个文件,其中每一行都是一个json 对象

{id_1: {"Samsung":{"popularity_score":0.4}, "long_battery":{"popularity_score": 0.8}}}
{id_2: {"Apple": {"popularity_score": 0.9},"UI":{"popularity_score":0.9}}}

我尝试使用 to_jsoncollect_list 函数,但我得到一个列表而不是嵌套的 json 对象。 这是一个大型分布式数据帧,因此不能转换为 pandas 或将其收集到一台机器中。

【问题讨论】:

    标签: python json apache-spark pyspark apache-spark-sql


    【解决方案1】:

    您需要为您的 JSON 创建一些地图类型:

    import pyspark.sql.functions as F
    
    df2 = df.groupBy('item_id').agg(
        F.map_from_entries(
            F.collect_list(
                F.struct('popular_tags', F.struct('popularity_score'))
            )
        ).alias('m')
    ).select(
        F.to_json(
            F.create_map('item_id', 'm')
        ).alias('col')
    )
    
    df2.show(truncate=False)
    +-------------------------------------------------------------------------------------+
    |col                                                                                  |
    +-------------------------------------------------------------------------------------+
    |{"id_2":{"Apple":{"popularity_score":0.9},"UI":{"popularity_score":0.9}}}            |
    |{"id_1":{"Samsung":{"popularity_score":0.4},"long battery":{"popularity_score":0.8}}}|
    +-------------------------------------------------------------------------------------+
    

    没有map_from_entries,你可能不得不依赖一些肮脏的黑客:

    df2 = df.groupBy('item_id').agg(
        F.collect_list(
            F.create_map('popular_tags', F.struct('popularity_score'))
        ).alias('m')
    ).select(
        F.regexp_replace(
            F.regexp_replace(
                F.to_json(F.create_map('item_id', 'm')),
                '(\\[|\\])', 
                ''
            ),
        '\\},\\{', 
        ','
        ).alias('col')
    )
    
    df2.show(truncate=False)
    +-------------------------------------------------------------------------------------+
    |col                                                                                  |
    +-------------------------------------------------------------------------------------+
    |{"id_2":{"Apple":{"popularity_score":0.9},"UI":{"popularity_score":0.9}}}            |
    |{"id_1":{"Samsung":{"popularity_score":0.4},"long battery":{"popularity_score":0.8}}}|
    +-------------------------------------------------------------------------------------+
    

    【讨论】:

    • 我有点明白你在做什么。谢谢,但我使用的 spark 2.2 没有这个 map_from_entries 函数。我可以使用 spark2.2 函数实现该功能吗?
    • @NG_21 ...这有点困难。我添加了一个 hacky 方法来做到这一点,不确定它是否适合你。
    • 我刚刚想出了一种方法,将您提到的解决方案与适用于 spark 2.2 的基于 UDF 的方法相结合。虽然更喜欢火花函数而不是 UDF 以获得性能。感谢您的帮助
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-11-18
    • 2021-12-21
    • 1970-01-01
    • 2019-07-08
    • 2021-08-26
    • 2020-04-02
    相关资源
    最近更新 更多