如何向ProcessPoolExecutor传递变量?遭遇TypeError错误
使用ProcessPoolExecutor.map()触发TypeError的解决方法
问题概述
尝试用ProcessPoolExecutor实现并行处理,使用map()向目标函数传递参数时抛出TypeError,提示float对象不可迭代。疑问:是否需要修改方法,让函数仅处理单行数据再循环传参?
错误信息
186 def _get_chunks(*iterables, chunksize): 187 """ Iterates over zip()ed iterables in chunks. """ --> 188 it = zip(*iterables) 189 while True: 190 chunk = tuple(itertools.islice(it, chunksize)) TypeError: 'float' object is not iterable
原代码
from concurrent.futures import ProcessPoolExecutor import pandas as pd import difflib import numpy as np with ProcessPoolExecutor() as executor: df = executor.map(drop_dupplicate_rows, df,'title',0.7,chunksize=100) def drop_dupplicate_rows(_df,_column,_cuttoff): ''' 基于difflib库去除重复行 参数 ---------- df: DataFrame 输入的数据集 _column: string 用于检测重复的列名 _cuttoff: float 判断重复的相似度阈值 返回值 ---------- 去重后的DataFrame ''' for k,row in enumerate(_df[f'{_column}']): repetition_list=list() if not np.where(_df[f'{_column}']==''): repetition_list=difflib.get_close_matches(row,_df[f'{_column}'],cutoff=_cuttoff) if len(repetition_list)>=2: print(f'row:{k},link: {_df.at[k,"url"]} due to: {_column} repetition_list: ',repetition_list) _df.drop(index=k,inplace=True) _df.reset_index(drop=True, inplace=True) return _df
问题原因
ProcessPoolExecutor.map()的核心规则是:第一个参数为目标函数,后续参数必须是可迭代对象。它会从每个可迭代对象中依次取元素,打包成参数组传给目标函数。
你当前直接传入了单个DataFrame、字符串'title'、浮点数0.7,这些都是非可迭代的单个对象。当内部执行zip(*iterables)时,会尝试将0.7作为可迭代对象处理,而float无法迭代,因此触发TypeError。
另外,原函数设计本身不适合并行:函数直接处理整个DataFrame,且包含inplace修改,进程间内存隔离的特性会导致数据同步问题,同时这种全局处理的逻辑也发挥不出并行的优势。
解决方案
方案1:调整为Chunk级并行处理
将大DataFrame拆分为多个小Chunk,每个Chunk由单独进程处理,最后合并结果。此方法适合处理大型数据集,注意:该方案仅处理Chunk内部的重复,跨Chunk的重复需要额外处理。
修改后的代码:
from concurrent.futures import ProcessPoolExecutor import pandas as pd import difflib import numpy as np from functools import partial def drop_dupplicate_rows(chunk_df, column, cutoff): '''按指定列,基于difflib去重单个DataFrame Chunk''' # 修正原空值判断逻辑 for k, row in enumerate(chunk_df[column]): if pd.isna(row) or row == '': continue # 在当前Chunk内查找匹配项 repetition_list = difflib.get_close_matches(row, chunk_df[column], cutoff=cutoff) if len(repetition_list) >= 2: print(f'row:{k}, link: {chunk_df.at[k,"url"]} due to: {column} repetition_list: {repetition_list}') chunk_df.drop(index=k, inplace=True) chunk_df.reset_index(drop=True, inplace=True) return chunk_df # 将DataFrame拆分为指定大小的Chunk def split_df(df, chunk_size=100): return [df.iloc[i:i+chunk_size] for i in range(0, len(df), chunk_size)] if __name__ == '__main__': # 替换为你的DataFrame加载逻辑 df = pd.read_csv("your_data_source.csv") df_chunks = split_df(df, chunk_size=100) # 用partial固定column和cutoff参数,避免map时传递多个可迭代对象 processed_func = partial(drop_dupplicate_rows, column='title', cutoff=0.7) with ProcessPoolExecutor() as executor: # 用map处理所有Chunk,转换为列表后合并 processed_chunks = list(executor.map(processed_func, df_chunks)) cleaned_df = pd.concat(processed_chunks, ignore_index=True)
方案2:优化单进程逻辑(无需并行)
原函数的循环效率极低,且空值判断逻辑有误。如果数据集规模不大,直接优化单进程逻辑比并行更高效:
import pandas as pd from fuzzywuzzy import process def fuzzy_drop_duplicates(df, column, cutoff=0.7): '''基于模糊匹配的全局去重,效率远高于循环单条处理''' # 过滤空值 non_null_df = df[df[column].notna() & (df[column] != '')].reset_index(drop=True) values = non_null_df[column].tolist() seen = set() keep_indices = [] for idx, val in enumerate(values): if val in seen: continue # 批量获取相似度高于cutoff的所有匹配项 matches = process.extract(val, values, limit=None, score_cutoff=int(cutoff * 100)) seen.update([match[0] for match in matches]) keep_indices.append(idx) # 返回去重后的DataFrame return non_null_df.iloc[keep_indices].reset_index(drop=True) # 使用示例 cleaned_df = fuzzy_drop_duplicates(your_df, column='title', cutoff=0.7)
关键注意事项
- 进程池处理DataFrame时,每个进程会复制Chunk数据,内存占用会提升,需根据机器配置调整Chunk大小。
- 原代码中
if not np.where(_df[f'{_column}']==''):逻辑错误,np.where返回索引数组,not判断永远为False,需替换为直接判断当前行是否为空/NaN。 - 若需要全局去重,Chunk并行方案会遗漏跨Chunk的重复,可先对整个列做相似度分组,再分配到进程处理,或先用单进程做一次全局预处理。
内容的提问来源于stack exchange,提问作者Mostafa Bouzari
相关产品推荐
相关产品推荐

