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

PySpark读取多CSV文件,仅处理数据量超指定阈值的文件

PySpark按文件阈值筛选读取CSV实现方案

问题根源

你现有代码的两个核心错误:

  • sc.textFile读取目录时会将所有CSV文件的行合并为同一个RDD,完全丢失行所属的文件关联信息,无法统计单个文件的行数/大小
  • x.split(header)逻辑错误:该操作是将每一行内容用表头完整字符串做分割,绝大多数行不包含表头,分割后长度最多为2,所以筛选长度>5的结果必然为空

正确实现方案

方案1:按单个文件行数筛选(和参考代码逻辑对齐)

参考Scala/Python示例的筛选逻辑是过滤行数不足阈值的股票数据文件,用wholeTextFilesAPI实现按文件粒度处理:

# 读取目录下所有文件,RDD元素结构为 (文件路径, 文件全部内容字符串)
stock_files_rdd = sc.wholeTextFiles("./data/stocks/")
# 定义行数阈值,和参考代码一致:5年每年260个交易日加10天冗余
ROW_THRESHOLD = 260 * 5 + 10

def parse_valid_file(content):
    # 拆分文件内容为行,跳过表头,过滤空行
    lines = content.strip().split("\n")
    return [line.strip() for line in lines[1:] if line.strip()]

# 先过滤行数达标的文件,再展开所有有效行
raw_stocks_rdd = stock_files_rdd.filter(lambda file_item: len(file_item[1].strip().split("\n")) - 1 >= ROW_THRESHOLD) \
                                .flatMap(lambda file_item: parse_valid_file(file_item[1]))

# 测试输出
print(raw_stocks_rdd.take(5))

方案2:按单个文件字节大小筛选

如果需要按文件占用存储空间大小筛选,先本地过滤符合要求的文件路径再读取即可:

import os

# 定义文件大小阈值,示例为1MB,可自行调整
SIZE_THRESHOLD = 1024 * 1024
stock_dir_path = "./data/stocks/"

# 筛选符合大小要求的CSV文件路径
valid_csv_paths = []
for fname in os.listdir(stock_dir_path):
    fpath = os.path.join(stock_dir_path, fname)
    if os.path.isfile(fpath) and fname.endswith(".csv") and os.path.getsize(fpath) >= SIZE_THRESHOLD:
        valid_csv_paths.append(fpath)

# 仅读取符合要求的文件
stocks_rdd = sc.textFile(",".join(valid_csv_paths))
# 跳过表头
header = stocks_rdd.first()
stocks_rdd = stocks_rdd.filter(lambda line: line != header)

# 测试输出
print(stocks_rdd.take(5))

内容的提问来源于stack exchange,提问作者starlight

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 07:45:02