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

如何将生成器流式传输至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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 15:17:26