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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 01:18:01