使用ThreadPoolExecutor和map多线程追加列表元素的异常问题
问题描述
编写了生成用户名、单线程生成CSV的函数,以及基于ThreadPoolExecutor和map的多线程生成CSV函数。已知列表append方法线程安全,但生成的CSV出现异常格式,例如:
[['EMP_ID', 'username', [234687, '234687Oregon696']]]
预期应为多个独立的用户数据行,而非表头行中嵌套用户数据。
相关代码
基础函数与单线程实现
import random import csv from concurrent.futures import ThreadPoolExecutor from concurrent.futures import ProcessPoolExecutor import pandas as pd def generate_username(id, job_location): number = "{:03d}".format(random.randrange(1, 999)) return "".join([id, job_location.strip(), str(number)]) def append_to_list(l, idx, job_location): l.append([idx, generate_username(str(idx), job_location)]) def generate_csv(filepath, df): rows = [["EMP_ID", "username"]] ids, locations = df.EMP_ID, df["Job Location"] for idx, location in zip(ids, locations): rows.append([idx, generate_username(str(idx), location)]) with open(filepath, 'w') as file: writer = csv.writer(file) writer.writerows(rows)
多线程实现(存在问题)
def generate_csv_threads(filepath, df, n): rows = [["EMP_ID", "username"]] ids, locations = df.EMP_ID, df["Job Location"] with ThreadPoolExecutor(max_workers=n) as executor: executor.map(append_to_list, rows, ids, locations) executor.shutdown(wait=True) with open(filepath, 'w') as file: writer = csv.writer(file) writer.writerows(rows)
问题原因
核心问题出在executor.map的参数传递逻辑上:
map会将第二个及以后的可迭代对象的元素逐个对应传给第一个函数的参数- 当前代码中,
rows是一个包含表头的列表([["EMP_ID", "username"]]),map会迭代rows的元素(即表头列表["EMP_ID", "username"]),将其作为append_to_list的第一个参数l - 而非把整个
rows列表作为l传入! - 最终导致所有用户数据都被append到表头列表内部,而非
rows列表中,从而出现表头行嵌套用户数据的异常格式。
修复方案
需要将rows作为固定参数绑定到append_to_list,确保每个线程都操作整个rows列表而非表头元素。以下是两种可行方案:
方案一:使用functools.partial绑定参数
from functools import partial def generate_csv_threads(filepath, df, n): rows = [["EMP_ID", "username"]] ids, locations = df.EMP_ID, df["Job Location"] # 绑定rows作为append_to_list的第一个参数 bound_append = partial(append_to_list, rows) with ThreadPoolExecutor(max_workers=n) as executor: # 仅传入idx和location两个参数 executor.map(bound_append, ids, locations) with open(filepath, 'w') as file: writer = csv.writer(file) writer.writerows(rows)
方案二:使用lambda表达式包装调用
def generate_csv_threads(filepath, df, n): rows = [["EMP_ID", "username"]] ids, locations = df.EMP_ID, df["Job Location"] with ThreadPoolExecutor(max_workers=n) as executor: executor.map(lambda idx, loc: append_to_list(rows, idx, loc), ids, locations) with open(filepath, 'w') as file: writer = csv.writer(file) writer.writerows(rows)
内容的提问来源于stack exchange,提问作者Norhther
相关产品推荐
相关产品推荐

