【问题标题】:How do I reduce a spark dataframe to a maximum amount of rows for each value in a column?如何将火花数据框减少到列中每个值的最大行数?
【发布时间】:2020-04-27 15:29:04
【问题描述】:

我需要减少数据名并将其导出到镶木地板。我需要确保我有前任。一列中的每个值对应 10000 行。

我正在使用的数据框如下所示:

+-------------+-------------------+
|         Make|              Model|
+-------------+-------------------+
|      PONTIAC|           GRAND AM|
|        BUICK|            CENTURY|
|        LEXUS|             IS 300|
|MERCEDES-BENZ|           SL-CLASS|
|      PONTIAC|           GRAND AM|
|       TOYOTA|              PRIUS|
|   MITSUBISHI|      MONTERO SPORT|
|MERCEDES-BENZ|          SLK-CLASS|
|       TOYOTA|              CAMRY|
|         JEEP|           WRANGLER|
|    CHEVROLET|     SILVERADO 1500|
|       TOYOTA|             AVALON|
|         FORD|             RANGER|
|MERCEDES-BENZ|            C-CLASS|
|       TOYOTA|             TUNDRA|
|         FORD|EXPLORER SPORT TRAC|
|    CHEVROLET|           COLORADO|
|   MITSUBISHI|            MONTERO|
|        DODGE|      GRAND CARAVAN|
+-------------+-------------------+

我需要为每个模型返回最多 10,000 行:

+--------------------+-------+
|               Model|  count|
+--------------------+-------+
|                 MDX|1658647|
|               ASTRO| 682657|
|           ENTOURAGE|  72622|
|             ES 300H|  80712|
|            6 SERIES| 145252|
|           GRAN FURY|   9719|
|RANGE ROVER EVOQU...|   4290|
|        LEGACY WAGON|   2070|
|        LEGACY SEDAN|    104|
|  DAKOTA CHASSIS CAB|      8|
|              CAMARO|2028678|
|                  XT|  10009|
|             DYNASTY| 171776|
|                 944|  43044|
|         F430 SPIDER|    506|
|FLEETWOOD SEVENTY...|      6|
|         MONTE CARLO|1040806|
|             LIBERTY|2415456|
|            ESCALADE| 798832|
| SIERRA 3500 CLASSIC|   9541|
+--------------------+-------+

This question 不一样,因为正如其他人在下面建议的那样,它只检索值大于其他值的行。我想要for each value in df['Model']: limit rows for that value(model) to 10,000 if there are 10,000 or more rows(显然是伪代码)。换言之,如果超过 10,000 行,则删除其余行,否则保留所有行。

【问题讨论】:

  • 将 distinct() 不做你需要的......我不认为我理解你的问题。你能提供更多的代码和数据吗..
  • 窗口函数是一种方式。您还可以按品牌和模型进行分区,然后使用索引映射分区并传入一个函数,该函数返回索引小于 1000 的所有行。
  • @Ravaal 请更详细地解释为什么发布的副本没有回答您的问题。最好的方法是创建一个带有少量模型的minimal reproducible example,并将阈值设为一个较小的数字(如 3)以准确显示所需的输出。
  • @pault 我已经尝试解释了很多次了。我不明白这有什么难理解的。如果每个值的行数超过 10K,则只返回 10K ROWS FOR THAT VALUE。如果行数少于 10K,则返回所有行。这会给我一个小得多的样本,我可以使用它。不幸的是,我无法进入我的实例并使用代码。这对我的公司来说非常昂贵......

标签: apache-spark pyspark pyspark-dataframes


【解决方案1】:

我想您应该将row_numberwindoworderBypartitionBy 放在一起查询结果,然后您可以根据您的限制进行过滤。例如,获得随机洗牌并将样本限制为每个值 10,000 行,如下所示:

from pyspark.sql import functions as F
from pyspark.sql.window import Window

window = Window.partitionBy(df['Model']).orderBy(F.rand())
df = df.select(F.col('*'), 
               F.row_number().over(window).alias('row_number')) \
               .where(F.col('row_number') <= 10000)

【讨论】:

  • 恭喜 slmn!你的帖子是最接近的。我进行了重大修改,但你还是赢得了赏金!
【解决方案2】:

如果我理解您的问题,您希望抽样几行(例如 10000),但这些记录的计数应该大于 10000。如果我理解您的问题,这就是答案:

df = df.groupBy('Make', 'Model').agg(count(lit(1)).alias('count'))
df = df.filter(df['count']>10000).select('Model','count')
df.write.parquet('output.parquet')

【讨论】:

  • @Ravaal 我已经添加了限制,因为我认为您需要按照您所说的“采样”,而不是获取所有记录。好的不错!那么问题已经回答了吗?
  • 我不得不收回答案。它没有用。有关实际数据框和说明,请参阅我的编辑。
  • @Ravaal 我再次编辑了我的答案,看看并告诉我它是否有效,如果不是什么问题。
  • 完美,所以@Ravaal 问题得到解答了吗?
  • count() 未定义。我试过from pyspark.sql.functions import count,但是count需要一个参数,所以我将'Model'作为参数传递给count。然后我输入df = df.filter(df['count']&gt;10000).select('Model','count') 并跟进df.count() 并得到0 个结果。很遗憾,我们在这里没有答案。
【解决方案3】:

简单地做

import pyspark.sql.functions as F

df = df.groupBy("Model").agg(F.count(F.lit(1)).alias("Count"))
df = df.filter(df["Count"] < 10000).select("Model", "Count")

df.write.parquet("data.parquet")

【讨论】:

  • 那行不通。我需要模型列中每个值的样本...
  • 我添加了distinct() tho,这将满足您的需求
  • 这将是理想的。
  • 它没有成功。数据帧从超过 3 亿条记录变为只有 1699 条记录。不过这个主意不错。
  • 这样就可以了。这意味着您有 1699 个唯一值。所以你不能有比唯一值更多的行
【解决方案4】:

我将稍微修改给定的问题,以便在此处可视化,方法是将每个不同值的最大行数减少到 2 行(而不是 10,000 行)。

示例数据框:

df = spark.createDataFrame(
  [('PONTIAC', 'GRAND AM'), ('BUICK', 'CENTURY'), ('LEXUS', 'IS 300'), ('MERCEDES-BENZ', 'SL-CLASS'), ('PONTIAC', 'GRAND AM'), ('TOYOTA', 'PRIUS'), ('MITSUBISHI', 'MONTERO SPORT'), ('MERCEDES-BENZ', 'SLK-CLASS'), ('TOYOTA', 'CAMRY'), ('JEEP', 'WRANGLER'), ('MERCEDES-BENZ', 'SL-CLASS'), ('PONTIAC', 'GRAND AM'), ('TOYOTA', 'PRIUS'), ('MITSUBISHI', 'MONTERO SPORT'), ('MERCEDES-BENZ', 'SLK-CLASS'), ('TOYOTA', 'CAMRY'), ('JEEP', 'WRANGLER'), ('CHEVROLET', 'SILVERADO 1500'), ('TOYOTA', 'AVALON'), ('FORD', 'RANGER'), ('MERCEDES-BENZ', 'C-CLASS'), ('TOYOTA', 'TUNDRA'), ('TOYOTA', 'PRIUS'), ('MITSUBISHI', 'MONTERO SPORT'), ('MERCEDES-BENZ', 'SLK-CLASS'), ('TOYOTA', 'CAMRY'), ('JEEP', 'WRANGLER'), ('CHEVROLET', 'SILVERADO 1500'), ('TOYOTA', 'AVALON'), ('FORD', 'RANGER'), ('MERCEDES-BENZ', 'C-CLASS'), ('TOYOTA', 'TUNDRA'), ('FORD', 'EXPLORER SPORT TRAC'), ('CHEVROLET', 'COLORADO'), ('MITSUBISHI', 'MONTERO'), ('DODGE', 'GRAND CARAVAN')],
  ['Make', 'Model']
)

让我们做一个行数:

df.groupby('Model').count().collect()

+-------------------+-----+
|              Model|count|
+-------------------+-----+
|             AVALON|    2|
|            CENTURY|    1|
|             TUNDRA|    2|
|           WRANGLER|    3|
|           GRAND AM|    3|
|EXPLORER SPORT TRAC|    1|
|            C-CLASS|    2|
|      MONTERO SPORT|    3|
|              CAMRY|    3|
|      GRAND CARAVAN|    1|
|     SILVERADO 1500|    2|
|              PRIUS|    3|
|            MONTERO|    1|
|           COLORADO|    1|
|             RANGER|    2|
|          SLK-CLASS|    3|
|           SL-CLASS|    2|
|             IS 300|    1|
+-------------------+-----+

如果我正确理解您的问题,您可以通过Model 为每一行分配一个行号:

from pyspark.sql import Window
from pyspark.sql.functions import row_number, desc

win_1 = Window.partitionBy('Model').orderBy(desc('Make'))
df = df.withColumn('row_num', row_number().over(win_1))

然后将行过滤到row_num &lt;= 2:

df = df.filter(df.row_num <= 2).select('Make', 'Model')

总共应该有2+1+2+2+2+1+2+2+2+1+2+2+1+1+2+2+2+1 = 30行

最终结果:

+-------------+-------------------+
|         Make|              Model|
+-------------+-------------------+
|       TOYOTA|             AVALON|
|       TOYOTA|             AVALON|
|        BUICK|            CENTURY|
|       TOYOTA|             TUNDRA|
|       TOYOTA|             TUNDRA|
|         JEEP|           WRANGLER|
|         JEEP|           WRANGLER|
|      PONTIAC|           GRAND AM|
|      PONTIAC|           GRAND AM|
|         FORD|EXPLORER SPORT TRAC|
|MERCEDES-BENZ|            C-CLASS|
|MERCEDES-BENZ|            C-CLASS|
|   MITSUBISHI|      MONTERO SPORT|
|   MITSUBISHI|      MONTERO SPORT|
|       TOYOTA|              CAMRY|
|       TOYOTA|              CAMRY|
|        DODGE|      GRAND CARAVAN|
|    CHEVROLET|     SILVERADO 1500|
|    CHEVROLET|     SILVERADO 1500|
|       TOYOTA|              PRIUS|
|       TOYOTA|              PRIUS|
|   MITSUBISHI|            MONTERO|
|    CHEVROLET|           COLORADO|
|         FORD|             RANGER|
|         FORD|             RANGER|
|MERCEDES-BENZ|          SLK-CLASS|
|MERCEDES-BENZ|          SLK-CLASS|
|MERCEDES-BENZ|           SL-CLASS|
|MERCEDES-BENZ|           SL-CLASS|
|        LEXUS|             IS 300|
+-------------+-------------------+

【讨论】:

  • 请缩小该图像。
猜你喜欢
  • 2020-07-06
  • 2019-07-13
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-07-13
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多