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

为何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...
...

问题根源

  1. 函数参数未被使用:executor.map会把users_list的每个元素传给process_members函数,但函数里完全没用到传入的user_list参数,反而硬编码取全局users_list[0],导致所有进程都处理第一个元素。
  2. 全局列表的多进程竞争: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 14:23:21