【问题标题】:PySpark join dataframes and merge contents of specific columnsPySpark 加入数据框并合并特定列的内容
【发布时间】:2018-07-11 16:51:01
【问题描述】:

我的目标是合并列 id 上的两个数据框,并在包含我们可以称为 data 的 JSON 的另一列上执行一些复杂的合并。

假设我的 DataFrame df1 看起来像这样:

id | data
---------------------------------
42 | {'a_list':['foo'],'count':1}
43 | {'a_list':['scrog'],'count':0}

我有兴趣与相似但不同的 DataFrame df2 合并:

id | data
---------------------------------
42 | {'a_list':['bar'],'count':2}
44 | {'a_list':['baz'],'count':4}

我想要以下 DataFrame,加入和合并来自 JSON 数据的属性,其中 id 匹配,但保留 id 不匹配的行并保持 data 列原样:

id | data
---------------------------------------
42 | {'a_list':['foo','bar'],'count':3}  <-- where 'bar' is added to 'foo', and count is summed
43 | {'a_list':['scrog'],'count':1}
44 | {'a_list':['baz'],'count':4}

可以看出id 是 42,我必须将一些逻辑应用于 JSON 的合并方式。

我的下意识的想法是我想提供一个 lambda / udf 来合并 data 列,但不知道在加入期间如何考虑。

或者,我可以将 JSON 中的属性分成列,像这样,这可能是更好的方法吗?

df1:

id | a_list    | count
----------------------
42 | ['foo']   | 1
43 | ['scrog'] | 0

df2:

id | a_list   | count
---------------------
42 | ['bar']  | 2
44 | ['baz']  | 4

结果:

id | a_list         | count
---------------------------
42 | ['foo', 'bar'] | 3
43 | ['scrog']      | 0
44 | ['baz']        | 4

如果我走这条路,那么我将不得不将 a_listcount 列再次合并到 JSON 中的单个列 data 下,但是我可以把我的头换成一个相对简单的 map功能。

更新:扩展问题

更现实地说,我将在一个列表中拥有n 数量的 DataFrame,例如df_list = [df1, df2, df3],形状都一样。在 n 个 DataFrame 上执行这些相同操作的有效方法是什么?

更新到更新

不确定这是多么有效,或者是否有更火花的方式来做到这一点,但结合接受的答案,这似乎适用于问题更新:

for i in range(0, (len(validations) - 1)):  

    # set dfs
    df1 = validations[i]['df']
    df2 = validations[(i+1)]['df']

    # joins here...

    # update new_df
    new_df = df2

【问题讨论】:

标签: pyspark apache-spark-sql


【解决方案1】:

这是完成第二种方法的一种方法:

分解列表列,然后分解unionAll 两个DataFrame。下一组按“id”列并使用pyspark.sql.functions.collect_list()pyspark.sql.functions.sum()

import pyspark.sql.functions as f
new_df = df1.select("id", f.explode("a_list").alias("a_values"), "count")\
    .unionAll(df2.select("id", f.explode("a_list").alias("a_values"), "count"))\
    .groupBy("id")\
    .agg(f.collect_list("a_values").alias("a_list"), f.sum("count").alias("count"))

new_df.show(truncate=False)
#+---+----------+-----+
#|id |a_list    |count|
#+---+----------+-----+
#|43 |[scrog]   |0    |
#|44 |[baz]     |4    |
#|42 |[foo, bar]|3    |
#+---+----------+-----+

最后,您可以使用pyspark.sql.functions.struct()pyspark.sql.functions.to_json() 将此中间DataFrame 转换为您想要的结构:

new_df = new_df.select("id", f.to_json(f.struct("a_list", "count")).alias("data"))
new_df.show()
#+---+----------------------------------+
#|id |data                              |
#+---+----------------------------------+
#|43 |{"a_list":["scrog"],"count":0}    |
#|44 |{"a_list":["baz"],"count":4}      |
#|42 |{"a_list":["foo","bar"],"count":3}|
#+---+----------------------------------+

更新

如果您在 df_list 中有一个数据框列表,您可以执行以下操作:

from functools import reduce   # for python3
df_list = [df1, df2]
new_df = reduce(lambda a, b: a.unionAll(b), df_list)\
    .select("id", f.explode("a_list").alias("a_values"), "count")\
    .groupBy("id")\
    .agg(f.collect_list("a_values").alias("a_list"), f.sum("count").alias("count"))\
    .select("id", f.to_json(f.struct("a_list", "count")).alias("data"))

【讨论】:

  • 这绝对是正确的,就像我发布的问题的梦想一样。选择作为答案。但是,如果我可以更进一步,我实际上有 n 的 DataFrame 数量,您是否考虑扩展您的答案以解释如何遍历 DataFrame 列表并执行此操作?我已经更新了问题。
  • 不确定是否有效,但更新了问题似乎可行。
  • @ghukill 我提供了适用于数据框列表的更新。也不确定效率。
  • 完美运行,非常感谢。我将运行一些基准测试,看看我是否无法了解哪个更有效。
猜你喜欢
  • 2022-01-12
  • 2018-02-17
  • 1970-01-01
  • 1970-01-01
  • 2012-11-15
  • 1970-01-01
  • 1970-01-01
  • 2021-11-17
  • 1970-01-01
相关资源
最近更新 更多