为何executor.map未实现列表索引递增?并行处理问题排查
并行处理Graph API调用时重复获取同一TeamID的问题解决
问题背景
用for循环处理2000+次Graph API调用性能太差,改用concurrent.futures.ProcessPoolExecutor的executor.map实现并行后,所有进程都重复处理同一个TeamID,输出全是重复的ID,无法迭代处理列表中的不同项。原代码试图用array_value模拟索引递增,但完全没起作用。
原代码如下:
users_and_teams_list = [] team_display_name = [] length = len(users_list)-1 teams_list = [] #for i in range(length): def process_members(user_list): array_value = 0 TeamID = users_list[array_value]["id"] graph_api_endpoint = "https://graph.microsoft.com/v1.0/teams/"+TeamID+"/members" access_token = get_token_for_service_principal( client_id=client_id, client_secret=client_secret, tenantId=tenantId, scope=scope ) userData = get_data_using_token(userUrl, access_token) #Set headers for the API request headers = { "Authorization": "Bearer " + access_token, "Content-Type": "application/json", } response = requests.get(graph_api_endpoint, headers=headers) team_display_name.extend([users_list[array_value]["displayName"]]) if response.status_code == 200: teams_data = response.json().get("value", []) teams_list.extend(teams_data) users_and_teams_list.extend(team_display_name+teams_list) array_value = array_value + 1 # Check for nextLink for pagination while "@odata.nextLink" in response.json(): next_link = response.json()["@odata.nextLink"] # Make the next GET request to the nextLink for the next page response = requests.get(next_link, headers=headers) if response.status_code == 200: teams_data = response.json().get("value", []) teams_list.extend(teams_data) else: print("Error:", response.status_code) print(response.text) break else: print("Error:", response.status_code) print(response.text) print(TeamID + " was processed...") with concurrent.futures.ProcessPoolExecutor() as executor: executor.map(process_members, users_list)
实际输出全是重复的TeamID:
33089--75a3 was processed... 33089--75a3 was processed... ...
问题根源
- 函数参数未被使用:
executor.map会把users_list的每个元素传给process_members函数,但函数里完全没用到传入的user_list参数,反而硬编码取全局users_list[0],导致所有进程都处理第一个元素。 - 全局列表的多进程竞争:
users_and_teams_list、team_display_name、teams_list都是全局变量,多进程同时修改会导致数据混乱,而且array_value是函数内局部变量,每次调用都会重置为0,根本无法实现递增。
解决方案
修改思路
- 让
process_members直接接收单个用户对象,从对象里取TeamID和displayName,不需要依赖全局索引。 - 每个进程独立处理数据,返回结果给主进程汇总,避免全局共享状态。
- 保持原有的分页逻辑,确保每个Team的成员数据完整获取。
修改后的代码
import concurrent.futures import requests def process_members(user): # 从传入的user对象直接获取ID和displayName team_id = user["id"] team_display_name = user["displayName"] graph_api_endpoint = f"https://graph.microsoft.com/v1.0/teams/{team_id}/members" # 获取token(如果token有效期足够,建议在主进程获取一次后传入,避免重复调用) access_token = get_token_for_service_principal( client_id=client_id, client_secret=client_secret, tenantId=tenantId, scope=scope ) # 这里的userData看起来没用到,可以考虑移除 # userData = get_data_using_token(userUrl, access_token) headers = { "Authorization": f"Bearer {access_token}", "Content-Type": "application/json", } all_team_members = [] response = requests.get(graph_api_endpoint, headers=headers) if response.status_code == 200: all_team_members.extend(response.json().get("value", [])) # 处理分页 while "@odata.nextLink" in response.json(): next_link = response.json()["@odata.nextLink"] response = requests.get(next_link, headers=headers) if response.status_code == 200: all_team_members.extend(response.json().get("value", [])) else: print(f"Error fetching members for team {team_id}: {response.status_code}") print(response.text) break print(f"{team_id} was processed...") # 返回当前team的结果:displayName + 成员列表 return (team_display_name, all_team_members) else: print(f"Error fetching members for team {team_id}: {response.status_code}") print(response.text) return None # 主进程汇总结果 users_and_teams_list = [] with concurrent.futures.ProcessPoolExecutor() as executor: # 用map把users_list的每个元素传给process_members results = executor.map(process_members, users_list) # 遍历结果,合并到最终列表 for result in results: if result is not None: display_name, members = result users_and_teams_list.append(display_name) users_and_teams_list.extend(members)
额外优化建议
- Token复用:
get_token_for_service_principal的调用可以放在主进程执行一次,然后把token传给process_members函数,避免每个进程重复获取token,提升性能。 - 进程数控制:可以在
ProcessPoolExecutor里指定max_workers参数,比如ProcessPoolExecutor(max_workers=10),避免进程过多导致API限流或系统资源耗尽。
内容的提问来源于stack exchange,提问作者Jake
相关产品推荐
相关产品推荐

