如何在PySpark中选取每第50行数据?大DataFrame无法转Pandas
大规模DataFrame抽样解决方案(保持title_code平衡)
如果你用的是分布式DataFrame(比如Spark),完全不需要转Pandas就能实现需求,以下是两种针对性方案:
方案1:分组抽样(优先推荐,保证title_code平衡)
直接全局取每50条可能导致部分title_code的样本占比失衡,分组后在每个title_code内部抽样,能严格保证各分组的样本分布一致:
from pyspark.sql.window import Window from pyspark.sql.functions import row_number # 按title_code分组,组内可指定排序字段(这里用title_code示例) window = Window.partitionBy("title_code").orderBy("title_code") # 给每组内的记录添加行号 df_with_row = df.withColumn("row_num", row_number().over(window)) # 筛选每组内第50、100、150...条数据 sampled_df = df_with_row.filter("row_num % 50 == 0")
方案2:全局排序后抽样(匹配你原本的思路)
如果坚持要全局排序后取每50条,可通过全局窗口添加行号实现,注意这种方式可能无法保证title_code的平衡:
from pyspark.sql.window import Window from pyspark.sql.functions import row_number # 全局按title_code排序,添加全局行号 global_window = Window.orderBy("title_code") df_global_row = df.withColumn("global_row", row_number().over(global_window)) # 选取每第50条数据 sampled_df = df_global_row.filter("global_row % 50 == 0")
注意事项
- 若使用其他分布式框架(如Dask),核心逻辑一致:通过窗口/分区添加行号后筛选,避免转换为Pandas
- 全局排序对超大规模数据性能影响较大,优先选择方案1的分组处理,减少数据 shuffle 开销
- 如果每组数据量不足50条,方案1不会抽取该组的样本,可根据需求调整筛选条件(比如取每组的最后1条或按比例抽样)
内容的提问来源于stack exchange,提问作者user12545095
相关产品推荐
相关产品推荐

