Polars懒加载DataFrame显式类型转换报错,如何实现批量Parquet文件的合并转换?
我完全理解你的困惑——LazyFrame本就是为大数据场景设计的,结果遇到这种schema不一致的问题反而卡壳了。咱们先搞清楚问题出在哪,再给你两个靠谱的Lazy模式解决方案。
问题根源
你遇到的SchemaError本质是Lazy模式下批量扫描Parquet时,Polars会先尝试统一所有文件的Schema。比如某个文件里rain_rate是binary类型,另一个是f64,Polars在扫描阶段就会发现类型冲突,直接报错——而你的cast操作是在扫描之后才执行的,根本没机会处理这个前置的Schema不兼容问题。
而Eager模式没问题,是因为你每个文件单独读取、转换、再合并,相当于绕开了全局Schema校验的步骤。
方案1:逐个扫描+Lazy合并(和Eager逻辑对齐,大数据友好)
这个思路和你Eager版本的逻辑完全一致,但全程用LazyFrame执行,不会把所有数据加载到内存:
import glob import polars as pl schema = { 'station_id': pl.String, 'datetime_utc': pl.Datetime(time_unit='ns', time_zone='UTC'), 'rain_rate': pl.Float64, } # 逐个扫描每个文件,单独做类型转换 lazy_dfs = [] for file in glob.glob('../output/extraction/part_*.parquet'): lf = pl.scan_parquet(file) lf = lf.with_columns([ pl.col(name).cast(dtype, strict=False) for name, dtype in schema.items() ]) lazy_dfs.append(lf) # 合并所有LazyFrame,然后流式写入 pl.concat(lazy_dfs).sink_parquet('../output/extraction/merged.parquet')
这个方法的优势是和你已有的Eager代码逻辑完全匹配,而且全程是懒执行,即使文件数量多、数据量大,也不会占用过多内存——Polars会在sink_parquet的时候流式处理每个文件。
方案2:手动指定Schema+关闭严格校验(更简洁)
Polars的scan_parquet支持直接传入schema参数,强制指定读取时的列类型,再配合strict=False允许类型自动转换,一步到位解决问题:
import polars as pl schema = { 'station_id': pl.String, 'datetime_utc': pl.Datetime(time_unit='ns', time_zone='UTC'), 'rain_rate': pl.Float64, } lf = pl.scan_parquet( '../output/extraction/part_*.parquet', schema=schema, strict=False # 允许Polars自动尝试类型转换,忽略非致命的不兼容 ) lf.sink_parquet('../output/extraction/merged.parquet')
这个方案更简洁,不需要循环处理每个文件。schema参数告诉Polars“就按我指定的类型读”,strict=False让它遇到类型不一致时自动尝试转换(和你代码里的cast(strict=False)逻辑一致),完美绕开了扫描阶段的Schema冲突。
两种方案都能实现你要的Lazy流式处理,方案1更灵活(比如可以在每个文件上加额外的处理逻辑),方案2更简洁。根据你的场景选就行~
备注:内容来源于stack exchange,提问作者Droid

