【问题标题】:pyspark - create DataFrame Grouping columns in map type structurepyspark - 在地图类型结构中创建 DataFrame 分组列
【发布时间】:2018-01-13 21:22:21
【问题描述】:

我的DataFrame结构如下:

-------------------------
| Brand | type |  amount|
-------------------------
|  B   |   a  |   10   |
|  B   |   b  |   20   |
|  C   |   c  |   30   |
-------------------------

我想通过将 typeamount 分组到 type 的一列中来减少行数:Map 所以Brand 将是唯一的,MAP_type_AMOUNT 将有key,value 对应每个type amount 组合。

我认为 Spark.sql 可能有一些函数可以在这个过程中提供帮助,还是我必须将 RDD 作为 DataFrame 并“自己”转换为地图类型?

预期

   -------------------------
    | Brand | MAP_type_AMOUNT 
    -------------------------
    |  B    | {a: 10, b:20} |
    |  C    | {c: 30}       |
    -------------------------

【问题讨论】:

  • 我认为我们在数据框中没有类似 rdd 的 collectAsmap 的东西。更好的是你可以使用自己的转换函数。

标签: python sql dictionary pyspark spark-dataframe


【解决方案1】:

Prem's 的回答略有改进(抱歉我还不能评论)

使用func.create_map 而不是func.struct。见documentation

import pyspark.sql.functions as func
df = sc.parallelize([('B','a',10),('B','b',20),
('C','c',30)]).toDF(['Brand','Type','Amount'])

df_converted = df.groupBy("Brand").\
    agg(func.collect_list(func.create_map(func.col("Type"),
    func.col("Amount"))).alias("MAP_type_AMOUNT"))

print df_converted.collect()

输出:

[Row(Brand=u'B', MAP_type_AMOUNT=[{u'a': 10}, {u'b': 20}]),
 Row(Brand=u'C', MAP_type_AMOUNT=[{u'c': 30}])]

【讨论】:

  • 完美改进:)
  • @osbon123 我们如何做相反的事情?将地图类型转换为多列数据框?
  • 这不会生成字典,而是生成一个字典数组,每个字典只有一个键值对。
【解决方案2】:

你可以有类似下面的东西,但不完全是“地图”

import pyspark.sql.functions as func
df = sc.parallelize([('B','a',10),('B','b',20),('C','c',30)]).toDF(['Brand','Type','Amount'])

df_converted = df.groupBy("Brand").\
    agg(func.collect_list(func.struct(func.col("Type"), func.col("Amount"))).alias("MAP_type_AMOUNT"))
df_converted.show()

输出是:

+-----+----------------+
|Brand| MAP_type_AMOUNT|
+-----+----------------+
|    B|[[a,10], [b,20]]|
|    C|        [[c,30]]|
+-----+----------------+

【讨论】:

    【解决方案3】:

    同时使用collect_listmap_from_arrays可以实现这一点

    import pyspark.sql.functions as F
    
    df_converted = (
        df.groupBy('Brand')
        .agg(
            F.collect_list('type').alias('type'),
            F.collect_list('amount').alias('amount'),
        )
        .withColumn('MAP_type_AMOUNT', F.map_from_arrays('type', 'amount'))
        .drop('type', 'amount')
    )
    

    输出

    +-----+------------------+
    |Brand|   MAP_type_AMOUNT|
    +-----+------------------+
    |    C|         [c -> 30]|
    |    B|[b -> 20, a -> 10]|
    +-----+------------------+
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-01-01
      • 1970-01-01
      • 2023-03-19
      • 2021-04-10
      • 1970-01-01
      • 2016-05-02
      • 2018-12-19
      • 2022-01-20
      相关资源
      最近更新 更多