如何用Polars高效将大型CSV数据集按年月分区为Parquet?
大CSV按年月分区的Polars优化方案
问题根源
你当前代码的cache()并未生效——Polars的LazyFrame缓存的是查询计划而非实际数据。每次循环月份时,都会重新扫描整个CSV文件并执行年份过滤,导致重复开销,所以每个月处理时间几乎一致。
优化方法
1. 一次性流式分区写入(推荐)
只需扫描一次CSV,直接按年月分区生成Parquet文件,完全避免重复读取:
# 方式1:手动处理每个分区,生成独立文件 ( pl.scan_csv('some.csv', infer_schema_length=100000, null_values=['\\N']) .with_columns( pl.col('origin_datetime').dt.year().alias('year'), pl.col('origin_datetime').dt.month().alias('month') ) .partition_by('year', 'month') .collect(streaming=True) .map_batches(lambda batch: batch.write_parquet(f"/{batch['year'][0]}_{batch['month'][0]}.parquet")) ) # 方式2:利用Polars内置分区写入,生成规范目录结构 ( pl.scan_csv('some.csv', infer_schema_length=100000, null_values=['\\N']) .with_columns( pl.col('origin_datetime').dt.year().alias('year'), pl.col('origin_datetime').dt.month().alias('month') ) .collect(streaming=True) .write_parquet( root_path='./partitioned_data', partition_by=['year', 'month'] ) )
方式2会自动生成partitioned_data/year=2016/month=1/这类层级目录,每个目录下存储对应月份的Parquet文件,更符合大数据分区规范。
2. 先加载全量数据再分区(适合内存足够场景)
如果内存能容纳整个数据集,先一次性读取所有数据,再按年月过滤写入:
# 一次性读取并预处理全量数据 full_df = ( pl.scan_csv('some.csv', infer_schema_length=100000, null_values=['\\N']) .with_columns( pl.col('origin_datetime').dt.year().alias('year'), pl.col('origin_datetime').dt.month().alias('month') ) .collect(streaming=True) ) # 按年月循环写入,跳过空数据 for year in range(2016, 2019): year_subset = full_df.filter(pl.col('year') == year) for month in range(1, 13): month_subset = year_subset.filter(pl.col('month') == month) if not month_subset.is_empty(): month_subset.write_parquet(f"/{year}_{month}.parquet")
3. 额外细节优化
- 指定日期类型:若
origin_datetime自动解析异常,手动指定dtypes={'origin_datetime': pl.Datetime},减少类型转换开销。 - 手动定义Schema:提前明确所有列的类型,替代
infer_schema_length,进一步加快CSV读取速度。 - 内存友好参数:在
scan_csv中添加low_memory=True,适配内存紧张的运行环境。
内容的提问来源于stack exchange,提问作者romanzdk
相关产品推荐
相关产品推荐

