如何在Kedro中使用生成器?解决序列化报错问题
在Kedro中用生成器处理数据的解决办法
Kedro 默认使用的 MemoryDataSet 确实无法序列化生成器对象——因为生成器不能被 pickle,但这绝对不代表必须把所有数据都加载到内存里。你可以通过以下几种方式在 Kedro 中使用生成器:
1. 用流式数据集处理
Kedro 原生支持一些流式数据集,比如 TextFileDataSet,配置后可以按行读取数据,返回的是迭代器而非一次性加载全部内容。如果是自定义格式的数据,你可以继承 AbstractDataSet 实现自己的流式数据集:在 _load 方法里返回生成器,_save 方法里逐元素处理生成器的输出并写入存储。
2. 节点内部消化生成器,不返回给数据集
如果生成器只是节点内部处理数据的中间环节,完全可以在节点里直接遍历生成器完成计算,最后把计算后的非生成器结果返回给数据集。这样既享受到了生成器的内存优势,又避开了序列化的问题。
3. 用 LambdaDataSet 绕过序列化
如果一定要在节点间传递生成器,可以用 LambdaDataSet,它允许你自定义加载和保存逻辑,不需要序列化生成器。示例代码如下:
from kedro.io import LambdaDataSet def load_my_generator(): # 返回你的生成器 return (item for item in large_data_source) def save_my_generator(generator): # 自定义生成器的保存/处理逻辑 for item in generator: process_and_write(item) # 创建数据集实例 streaming_dataset = LambdaDataSet(load=load_my_generator, save=save_my_generator)
之后在 catalog.yml 里配置这个数据集,替换默认的 MemoryDataSet 即可。
4. 临时禁用结果保存(仅调试用)
如果只是调试阶段不想保存生成器的结果,可以在定义节点时把 outputs 设置为 None,这样 Kedro 就不会尝试序列化并保存返回的生成器:
def example_node(): return (item for item in large_data_source) # 在 pipeline 中定义节点 node( func=example_node, inputs=None, outputs=None, name="example_node" )
总结一下:Kedro 完全支持低内存的流式处理,只是默认的内存数据集不兼容生成器,通过自定义数据集、调整节点逻辑或者使用 LambdaDataSet 就能轻松解决这个问题。
内容的提问来源于stack exchange,提问作者ilja
相关产品推荐
相关产品推荐

