如何避免Dask中from_delayed为每个输入创建一个分区?
解决Dask from_delayed生成过多分区导致内存过载的问题
你当前的问题根源是:循环里每一条观测都单独生成一个delayed任务,调用from_delayed时每个delayed对象会被当作一个独立分区,最终分区数和观测数完全一致,直接导致任务图规模剧增、内存占用过载。
以下是几个可行的解决思路:
1. 批量处理观测,减少delayed任务数量
不要逐个处理单条数据,而是把数据分成若干批次,让每个delayed函数处理一整批观测,这样分区数就等于批次数量,能大幅压缩任务图规模。
修改后的代码示例:
@dask.delayed def match_batch_names(names_batch, countries_batch, cities_batch, name_series_orbis): matches = [] for name, country, city in zip(names_batch, countries_batch, cities_batch): best_match = process.extractOne(name, name_series_orbis)[0] matches.append({'Matched_Name': best_match}) return pd.DataFrame(matches) # 定义每批处理的数量,比如每批1000条 batch_size = 1000 lazy_results_names = [] # 把数据分成批次 for i in range(0, len(dask_dealscan_names_list), batch_size): names_batch = dask_dealscan_names_list[i:i+batch_size] countries_batch = dask_df_dealscan['Country'].iloc[i:i+batch_size].compute() cities_batch = dask_df_dealscan['City'].iloc[i:i+batch_size].compute() lazy_result = match_batch_names(names_batch, countries_batch, cities_batch, name_series_orbis) lazy_results_names.append(lazy_result) matched_names_dd = dd.from_delayed(lazy_results_names)
2. 使用Dask DataFrame的map_partitions替代循环+from_delayed
直接利用Dask DataFrame的map_partitions方法,对原DataFrame的每个分区应用匹配逻辑,天然继承原DataFrame的分区数,从源头避免生成过多分区。
修改后的代码示例:
def match_partition_names(partition_df, name_series_orbis): def match_single_row(row): best_match = process.extractOne(row['name_dealscan'], name_series_orbis)[0] return best_match # 假设原分区包含name_dealscan、Country、City列 partition_df['Matched_Name'] = partition_df.apply(match_single_row, axis=1) return partition_df[['Matched_Name']] # 假设dask_df_dealscan已经包含name_dealscan列(如果没有可以先合并) matched_names_dd = dask_df_dealscan.map_partitions( match_partition_names, name_series_orbis, meta=pd.DataFrame({'Matched_Name': pd.Series(dtype='object')}) )
3. 事后合并小分区(应急方案)
如果已经生成了过多分区的Dask DataFrame,可以用repartition或coalesce合并分区,但这只是应急手段,不如从源头控制分区数高效:
# coalesce会尽量合并现有分区,避免数据移动 matched_names_dd = matched_names_dd.coalesce(npartitions=100) # 按需设置目标分区数
内容的提问来源于stack exchange,提问作者Андрей Комиссаренко
相关产品推荐
相关产品推荐

