如何将生成器流式传输至Polars DataFrame并执行后续惰性计划?
问题描述
我有一个超长的生成器函数,希望使用Polars将其作为列处理。由于数据量庞大,我想以该生成器为源运行惰性流式处理,但尚未找到可行的实现方法(不确定是否可行)。
将生成器转为普通DataFrame再转换为惰性模式显然不可行,因为生成器会在调用collect()执行惰性计划前就被耗尽。LazyFrame初始化器也存在同样问题,它本质上只是上述操作的快捷方式。
请问是否存在无需先写入再扫描CSV的替代方案?
示例代码:
import polars as pl def Generator(): yield 1 yield 2 yield 3 generator = Generator() df = pl.DataFrame({"a": generator}).lazy() print(df) # naive plan... print([i for i in generator]) # [] generator2 = Generator() df = pl.LazyFrame({"a": generator2}) print(df) # naive plan... print([i for i in generator2]) # []
解决方案
用Polars的pl.scan_python就能解决这个问题,它支持从生成器惰性加载数据,不会提前耗尽生成器。需要注意的是,这个函数要求生成器输出的是行数据(比如单元素元组/列表,对应单列),具体用法如下:
直接适配生成器的写法
修改生成器让它返回元组形式的行:
import polars as pl def Generator(): yield (1,) # 单元素元组对应一列数据 yield (2,) yield (3,) # 创建惰性LazyFrame,指定schema避免类型推断消耗生成器 lazy_df = pl.scan_python(Generator(), schema={"a": pl.Int64}) # 此时生成器还没被消耗 print([i for i in Generator()]) # [(1,), (2,), (3,)] # 执行计算时才会读取生成器数据 result = lazy_df.collect() print(result) # shape: (3, 1) # ┌─────┐ # │ a │ # │ --- │ # │ i64 │ # ╞═════╡ # │ 1 │ # │ 2 │ # │ 3 │ # └─────┘
不修改原生成器的写法
如果不想改动原生成器,可以加一层包装函数,把单个值转成元组:
def wrap_generator(original_gen): for item in original_gen: yield (item,) lazy_df = pl.scan_python(wrap_generator(Generator()), schema={"a": pl.Int64}) # 后续操作和上面一致 result = lazy_df.collect() print(result)
为什么这个方法可行
pl.scan_python是Polars专门为Python迭代器/生成器设计的惰性数据源接口,它不会在初始化时就消耗生成器,只有当你调用collect()、fetch()等触发计算的方法时,才会逐步读取生成器的数据,完全符合流式处理超大生成器的需求。
内容的提问来源于stack exchange,提问作者AroneyS
相关产品推荐
相关产品推荐

