【发布时间】: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