You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 10:05:23