如何用Polars Lazy Frame处理多CSV文件的列不一致问题?
解决Polars扫描多CSV列数不一致的ShapeError问题
核心思路
先收集所有CSV文件的全部列名,再统一指定Schema扫描所有文件,缺失列自动填充Null,最后合并导出。
具体实现代码
import polars as pl from pathlib import Path # 1. 获取目标文件夹下所有CSV文件路径 csv_dir = "path" csv_files = list(Path(csv_dir).glob("*.csv")) # 2. 遍历所有文件,收集所有唯一列名 all_columns = set() for file in csv_files: # 仅读取表头,快速获取列名,避免加载整个文件 header_df = pl.read_csv(file, sep=";", n_rows=0) all_columns.update(header_df.columns) all_columns = list(all_columns) # 3. 构建统一Schema,所有列默认设为字符串类型(可根据实际需求调整类型) unified_schema = {col: pl.Utf8 for col in all_columns} # 4. 用统一Schema扫描所有CSV,缺失列自动填充Null combined_lf = pl.scan_csv( f"{csv_dir}/*.csv", sep=";", schema=unified_schema, infer_schema_length=0 ) # 5. 流式收集数据并写入Parquet combined_lf.collect(streaming=True).write_parquet("new_file.parquet")
备选简化方案(无需预收集列名)
如果不想提前遍历收集列名,也可以逐个扫描每个文件后统一选择所有列,再合并:
import polars as pl from pathlib import Path csv_files = list(Path("path").glob("*.csv")) # 先获取所有列名 all_columns = set() for f in csv_files: with pl.scan_csv(f, sep=";", n_rows=0) as lf: all_columns.update(lf.columns) all_columns = list(all_columns) # 逐个扫描并对齐列 lazy_frames = [] for f in csv_files: lf = pl.scan_csv(f, sep=";").select(all_columns) lazy_frames.append(lf) # 合并所有Lazy Frame combined_lf = pl.concat(lazy_frames, how="outer") combined_lf.collect(streaming=True).write_parquet("new_file.parquet")
注意事项
- 列类型可根据实际业务调整,比如把
pl.Utf8改成pl.Float64或pl.Int64,如果不确定类型,也可以用pl.Object(不推荐,会影响性能)。 - 对于超大量文件,预收集列名时用
pl.read_csv(..., n_rows=0)非常高效,因为只读取文件的表头部分。
内容的提问来源于stack exchange,提问作者miroslaavi
相关产品推荐
相关产品推荐

