【问题标题】:Creating a new dataframe with groupBy and filter使用 groupBy 和 filter 创建一个新的数据框
【发布时间】:2018-07-16 06:31:20
【问题描述】:

以下是我的dataframe

df = spark.createDataFrame([
  (0, 1),
  (0, 2),
  (0, 5),
  (1, 1),
  (1, 2),
  (1, 3),
  (1, 5),
  (2, 1),
  (2, 2)
], ["id", "product"])

我需要做一个groupByid 并收集所有项目,如下所示,但我需要检查产品数量,如果它小于2,那不应该在那里收集项目。

例如,产品 3 仅重复一次,即 3 的计数为 1,小于 2,因此它不应在以下数据帧中可用。看来我需要做两个groupBys:

预期输出:

+---+------------+
| id|       items|
+---+------------+
|  0|   [1, 2, 5]|
|  1|   [1, 2, 5]|
|  2|      [1, 2]|
+---+------------+

【问题讨论】:

  • 在哪个模块/库/包中找到 spark.createdataframe 函数?

标签: python python-3.x apache-spark pyspark apache-spark-sql


【解决方案1】:

我认为确实两个groupBy 是一个不错的解决方案,您可以在第一个groupBy 之后使用leftsemi 加入来过滤您的初始DataFrame。工作示例解决方案:

import pyspark.sql.functions as F

df = spark.createDataFrame([
(0, 1),
(0, 2),
(0, 5),
(1, 1),
(1, 2),
(1, 3),
(1, 5),
(2, 1),
(2, 2)
], ["id", "product"])

df = df\
.join(df.groupBy('product').count().filter(F.col('count')>=2),'product','leftsemi').distinct()\
.orderBy(F.col('id'),F.col('product'))\
.groupBy('id').agg(F.collect_list('product').alias('product'))

df.show()

orderBy 子句是可选的,仅当您关心结果中的顺序时。输出:

+---+---------+
| id|  product|
+---+---------+
|  0|[1, 2, 5]|
|  1|[1, 2, 5]|
|  2|   [1, 2]|
+---+---------+

希望这会有所帮助!

【讨论】:

    【解决方案2】:

    一种方法是使用Window 来获取每个产品的计数并将其用于filtergroupBy() 之前:

    import pyspark.sql.functions as f
    from pyspark.sql import Window
    
    df.withColumn('count', f.count('*').over(Window.partitionBy('product')))\
        .where('count > 1')\
        .groupBy('id')\
        .agg(f.collect_list('product').alias('items'))\
        .show()
    #+---+---------+
    #| id|    items|
    #+---+---------+
    #|  0|[5, 1, 2]|
    #|  1|[5, 1, 2]|
    #|  2|   [1, 2]|
    #+---+---------+
    

    不幸的是the HAVING statement does not exist in spark-sql

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-12-08
      • 2018-12-08
      • 2016-10-13
      • 1970-01-01
      • 2017-06-06
      • 2016-07-10
      • 1970-01-01
      相关资源
      最近更新 更多