【问题标题】:Updating json column using window cumulative via pyspark通过pyspark使用窗口累积更新json列
【发布时间】:2022-01-16 04:44:47
【问题描述】:

我一直在尝试通过 Pyspark、SparkSQL 和 Pandas 更新一系列 JSON blob,但没有成功。以下是数据的样子:

#+---+---------+------------------------------------------+
#|ID |Timestamp|Properties                                |
#+---+---------+------------------------------------------+
#|a  |7        |{"a1": 5, "a2": 8}                        |
#|b  |12       |{"b1": 36, "b2": "u", "b3": 17, "b8": "c"}|
#|a  |8        |{"a2": 4}                                 |
#|a  |10       |{"a3": "z", "a4": "t"}                    |
#|a  |5        |{"a1": 3, "a2": 12, "a4": "r"}            |
#|b  |20       |{"b2": "k", "b3": 9}                      |
#|b  |14       |{"b8": "y", "b3": 2}                      |
#+---+---------+------------------------------------------+

我想要一个查询,该查询将按ID 字段对行进行分区,并按Timestamp 字段对其进行排序。在此之后,Properties 字段将在每个分区中累积合并以创建一个新列New Props。所以输出是这样的:

#+---+---------+------------------------------------------+------------------------------------------+------+
#|ID |Timestamp|Properties                                |New_Props                                 |rownum|
#+---+---------+------------------------------------------+------------------------------------------+------+
#|a  |5        |{"a1": 3, "a2": 12, "a4": "r"}            |{"a1": 3, "a2": 12, "a4": "r"}            |1     |
#|a  |7        |{"a1": 5, "a2": 8}                        |{"a1": 5, "a2": 8, "a4": "r"}             |2     |
#|a  |8        |{"a2": 4}                                 |{"a1": 5, "a2": 4, "a4": "r"}             |3     |
#|a  |10       |{"a3": "z", "a4": "t"}                    |{"a1": 5, "a2": 4, "a3": "z", "a4": "t"}  |4     |
#|b  |12       |{"b1": 36, "b2": "u", "b3": 17, "b8": "c"}|{"b1": 36, "b2": "u", "b3": 17, "b8": "c"}|1     |
#|b  |14       |{"b8": "y", "b3": 2}                      |{"b1": 36, "b2": "u", "b3": 2, "b8": "y"} |2     |
#|b  |20       |{"b2": "k", "b3": 9}                      |{"b1": 36, "b2": "k", "b3": 9, "b8": "y"} |3     |
#+---+---------+------------------------------------------+------+------------------------------------------+

公式:从rownum2开始,获取上一行(rownum1)的New Props列值,并用当前行(rownum2)的Properties列的值更新.

我尝试使用 LAG 函数,但我无法使用我当前在函数本身内计算的列。

为了创建Next Props 列,我尝试了这个 CASE 语句,但它不起作用:

CASE
    WHEN rownum != 1 THEN concat(properties, LAG(next_props, 1) OVER (PARTITION BY contentid ORDER BY updateddatetime))
    ELSE next_props
END AS new_props

我一直在尝试不同的事情,但我被困住了。我可能可以使用 for 循环和 python dict.update() 函数来做到这一点,但我担心效率。任何帮助表示赞赏。

【问题讨论】:

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


    【解决方案1】:

    这是在数组和映射列上使用高阶函数的一种方法:

    1. 使用lag 为每一行获取前一个Properties,并将前一行和当前行Properties 转换为映射类型
    2. 在窗口上使用collect_list函数,得到前一行的累积数组Properties
    3. 将当前行Properties 添加到结果数组并使用aggregate 使用map_concat 连接内部映射。从您的示例来看,更新操作似乎只是添加新键,因此在 concat 之前,我们使用 map_filter 函数过滤已经存在的键
    4. 使用to_json从聚合映射中获取json字符串并删除中间列
    from pyspark.sql import functions as F, Window
    
    w = Window.partitionBy("ID").orderBy("Timestamp")
    
    df1 = df.withColumn("rownum", F.row_number().over(w)) \
        .withColumn("prev_prop_map", F.from_json(F.lag("Properties").over(w), "map<string,string>")) \
        .withColumn("current_prop_map", F.from_json("Properties", "map<string,string>")) \
        .withColumn("cumulative_prev_props", F.collect_list("prev_prop_map").over(w)) \
        .withColumn(
            "New_Props",
            F.to_json(F.aggregate(
                F.concat(F.array("current_prop_map"), F.reverse(F.col("cumulative_prev_props"))),
                F.expr("cast(map() as map<string,string>)"),
                lambda acc, x: F.map_concat(
                    acc,
                    F.map_filter(x, lambda k, _: ~F.array_contains(F.map_keys(acc), k))
                )
            ))
    ).drop("prev_prop_map", "current_prop_map", "cumulative_prev_props")
    
    
    df1.show(truncate=False)
    #+---+---------+------------------------------------------+------+---------------------------------------+
    #|ID |Timestamp|Properties                                |rownum|New_Props                              |
    #+---+---------+------------------------------------------+------+---------------------------------------+
    #|a  |5        |{"a1": 3, "a2": 12, "a4": "r"}            |1     |{"a1":"3","a2":"12","a4":"r"}          |
    #|a  |7        |{"a1": 5, "a2": 8}                        |2     |{"a1":"5","a2":"8","a4":"r"}           |
    #|a  |8        |{"a2": 4}                                 |3     |{"a2":"4","a1":"5","a4":"r"}           |
    #|a  |10       |{"a3": "z", "a4": "t"}                    |4     |{"a3":"z","a4":"t","a2":"4","a1":"5"}  |
    #|b  |12       |{"b1": 36, "b2": "u", "b3": 17, "b8": "c"}|1     |{"b1":"36","b2":"u","b3":"17","b8":"c"}|
    #|b  |14       |{"b8": "y", "b3": 2}                      |2     |{"b8":"y","b3":"2","b1":"36","b2":"u"} |
    #|b  |20       |{"b2": "k", "b3": 9}                      |3     |{"b2":"k","b3":"9","b8":"y","b1":"36"} |
    #+---+---------+------------------------------------------+------+---------------------------------------+
    

    如果您更喜欢使用 SQL 查询,这里是等效的 SparkSQL:

    WITH props AS (
        SELECT  *,
                row_number() over(partition by ID order by Timestamp) AS rownum,
                from_json(lag(Properties) over(partition by ID order by Timestamp), 'map<string,string>') AS prev_prop_map,
                from_json(Properties, 'map<string,string>') AS current_prop_map
        FROM    props_tb
    ),  cumulative_props AS (
        SELECT  *,
                collect_list(prev_prop_map) over(partition by ID order by Timestamp) AS cumulative_prev_props
        FROM    props 
    )
    
    SELECT  ID,
            Timestamp,
            Properties,
            aggregate(
                concat(array(current_prop_map), reverse(cumulative_prev_props)),
                cast(map() as map<string,string>),
                (acc, x) -> map_concat(acc, map_filter(x, (k,v) -> ! array_contains(map_keys(acc), k)))
            ) AS New_Props,
            rownum
    FROM    cumulative_props
    

    【讨论】:

    • 非常感谢您,这简直太棒了!我从未使用/听说过其中一些功能,例如collect_list()。为了获得我需要的输出,我只需将aggregate() 包装在to_json() 函数中,以再次将最终结果作为字符串。还有添加reverse() 函数的原因吗?我拿出来,结果还是正确的。
    猜你喜欢
    • 2020-11-17
    • 1970-01-01
    • 2021-03-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-09-24
    • 2020-03-10
    • 1970-01-01
    相关资源
    最近更新 更多