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

Polars懒加载DataFrame显式类型转换报错,如何实现批量Parquet文件的合并转换?

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 13:14:32