如何使用Dask实现千万级记录Python列表的并行过滤
问题根源
你现有实现存在两个核心问题:
- 单元素级的
dask.delayed调度开销极大,千万级数据下任务调度成本会远超过计算本身,生产环境完全不可用 - 自定义判断函数不满足条件时无显式返回,默认返回
None,才需要后续额外做串行过滤
最优方案:使用Dask Bag原生filter接口
Dask Bag是专门用来处理半结构化序列数据(比如字典列表)的组件,原生支持并行filter操作,不需要额外处理None,且按分区批量执行,调度成本极低。
你不需要修改原本的判断函数,直接用以下代码即可:
import dask.bag as db # 直接复用你最初写的返回布尔值的判断函数即可,不需要修改 def item_is_valid(item): item_is_valid = True if item['type'] == 'cat': item_is_valid = False elif item['weight'] > 20: item_is_valid = False # 其他判断条件保持不变 return item_is_valid # 将原始列表转为Dask Bag,npartitions可按你CPU核心数调整,千万级数据建议设置为50~200 items_bag = db.from_sequence(items, npartitions=8) # 直接调用filter方法,底层自动并行丢弃不满足条件的条目 filtered_bag = items_bag.filter(item_is_valid) # 执行计算,直接得到无None的过滤结果 items_filtered = filtered_bag.compute()
如果你的数据本身存储在磁盘的多个JSON/文本文件中,还可以直接用db.read_text()/db.read_json()直接读取文件,不需要先把全量数据加载到内存,进一步降低内存消耗。
备选方案:delayed批量分区处理
如果你想要更灵活地控制分区逻辑,也可以用按分区批量delayed的方式实现,同样不会产生多余的None:
from dask import delayed, compute # 将原始大列表拆分为固定大小的子块,单块大小建议1万~10万条,可根据实际调整 def split_chunks(lst, chunk_size=50000): return [lst[i:i+chunk_size] for i in range(0, len(lst), chunk_size)] item_chunks = split_chunks(items) @delayed def filter_single_chunk(chunk): # 单块内直接做本地过滤,返回的就是不含无效条目的子列表 return [item for item in chunk if item_is_valid(item)] # 并行处理所有块 delayed_chunks = [filter_single_chunk(chunk) for chunk in item_chunks] # 合并所有块的结果得到最终过滤列表 items_filtered = sum(compute(*delayed_chunks), [])
内容的提问来源于stack exchange,提问作者Victor
相关产品推荐
相关产品推荐

