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
相关产品推荐
相关产品推荐

