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

Python中批量并行运行For循环求助:Amazon Connect用户处理代码并行化未达预期

解决Amazon Connect用户批量处理的并行化问题

看起来你在处理大量Amazon Connect用户时遇到了性能瓶颈,想通过并行化加速循环但没搞定——作为新手这很正常,咱们一步步来解决这个问题。

先说说你现有代码里的几个关键问题:

  • 你在lambda_handler里递归调用了pool.starmap(lambda_handler, ...),这完全不对,相当于把整个handler又丢给进程池跑,根本没并行处理单个用户的逻辑
  • 原代码里有未定义的变量(比如routing_profile_name),还有security_profile_name_list是空的,这些都会导致报错
  • Boto3客户端不是进程安全的,直接在主进程创建客户端给子进程用可能会出问题
  • Lambda的执行环境有资源限制,进程池的大小要匹配你的Lambda配置(比如1vCPU的话别开超过4个进程)

接下来是修正后的完整方案:

步骤1:抽离单个用户的处理逻辑

把循环里处理单个用户的代码拆成独立函数,每个进程负责处理一个用户,并且在函数内部创建自己的Boto3客户端(避免进程安全问题)。

步骤2:正确使用进程池并行处理

在lambda_handler里获取用户列表后,用进程池来批量处理每个用户,最后收集结果。

完整修正代码

import json
import boto3
import os
import botocore
from datetime import datetime
import logging
import urllib3
import multiprocessing as mp
from botocore.client import Config

urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)

# 全局配置
config = Config(
    connect_timeout=7200,
    retries={'max_attempts': 5},
    proxies={
        'http': os.environ.get('HTTP_PROXY'),
        'https': os.environ.get('HTTPS_PROXY'),
    }
)
LOGGER = logging.getLogger()
LOGGER.setLevel(logging.INFO)

# 环境变量
instance_id = os.environ.get('instance_id')
region_name = os.environ.get('Region')
s3_bucket = os.environ.get('s3_bucket')
env = os.environ.get('export_environment')
MaxResults = 4

def getUserList(instance_id: str, region_name: str):
    """获取所有Amazon Connect用户列表"""
    client = boto3.client('connect', region_name=region_name, verify=False, config=config)
    users = []
    # 初始请求
    response = client.list_users(InstanceId=instance_id, MaxResults=MaxResults)
    users.extend(response['UserSummaryList'])
    # 分页获取剩余用户
    while "NextToken" in response:
        response = client.list_users(
            InstanceId=instance_id,
            MaxResults=MaxResults,
            NextToken=response["NextToken"]
        )
        users.extend(response['UserSummaryList'])
    return users

def process_user(user_data):
    """处理单个用户的逻辑,每个进程会调用这个函数"""
    user, instance_id, region_name = user_data
    # 每个进程创建自己的Boto3客户端
    client = boto3.client('connect', region_name=region_name, verify=False, config=config)
    
    user_id = user['Id']
    user_name = user['Username']
    
    # 获取用户详情
    try:
        response = client.describe_user(InstanceId=instance_id, UserId=user_id)
        user_details = response['User']
    except botocore.exceptions.ClientError as e:
        LOGGER.error(f"Failed to describe user {user_id}: {e}")
        return None
    
    # 提取基础信息
    identity_info = user_details['IdentityInfo']
    phone_config = user_details['PhoneConfig']
    
    user_firstname = identity_info.get('FirstName', '')
    user_lastname = identity_info.get('LastName', '')
    user_username = user_details['Username']
    user_routing_profile_id = user_details['RoutingProfileId']
    user_security_profile_ids = user_details['SecurityProfileIds']
    user_phonetype = phone_config.get('PhoneType', '')
    user_desk_number = phone_config.get('DeskPhoneNumber', '')
    user_autoaccept = phone_config.get('AutoAccept', False)
    user_acw_timeout = phone_config.get('AfterContactWorkTimeLimit', 0)
    user_hierarchy_group_id = user_details.get('HierarchyGroupId', '')
    
    # 获取路由配置名称
    routing_profile_name = ''
    try:
        rp_response = client.describe_routing_profile(
            InstanceId=instance_id,
            RoutingProfileId=user_routing_profile_id
        )
        routing_profile_name = rp_response['RoutingProfile']['Name']
    except botocore.exceptions.ClientError as e:
        LOGGER.warning(f"Failed to get routing profile for user {user_id}: {e}")
    
    # 获取安全配置名称列表
    security_profile_name_list = []
    for sp_id in user_security_profile_ids:
        try:
            sp_response = client.describe_security_profile(
                InstanceId=instance_id,
                SecurityProfileId=sp_id
            )
            security_profile_name_list.append(sp_response['SecurityProfile']['Name'])
        except botocore.exceptions.ClientError as e:
            LOGGER.warning(f"Failed to get security profile {sp_id} for user {user_id}: {e}")
    security_profile_names = ','.join(security_profile_name_list)
    
    # 获取层级组名称
    hierarchy_group_name = ''
    if user_hierarchy_group_id:
        try:
            hg_response = client.describe_hierarchy_group(
                InstanceId=instance_id,
                HierarchyGroupId=user_hierarchy_group_id
            )
            hierarchy_group_name = hg_response['HierarchyGroup']['Name']
        except botocore.exceptions.ClientError as e:
            LOGGER.warning(f"Failed to get hierarchy group {user_hierarchy_group_id} for user {user_id}: {e}")
    
    # 拼接用户信息行
    user_info = f"{user_firstname},{user_lastname},{user_username},{routing_profile_name},{security_profile_names},{user_phonetype},{user_desk_number},{str(user_autoaccept)},{str(user_acw_timeout)},{hierarchy_group_name}\r\n"
    return user_info

def lambda_handler(event, context):
    # 获取所有用户列表
    users = getUserList(instance_id, region_name)
    LOGGER.info(f"Found {len(users)} users to process")
    
    # 准备进程池:根据Lambda的CPU配置调整进程数,1vCPU建议2-4个
    pool_size = min(mp.cpu_count(), 4)
    with mp.Pool(processes=pool_size) as pool:
        # 构造每个进程需要的参数:(user, instance_id, region_name)
        user_tasks = [(user, instance_id, region_name) for user in users]
        # 并行处理所有用户
        results = pool.map(process_user, user_tasks)
    
    # 过滤掉处理失败的结果
    valid_results = [res for res in results if res is not None]
    LOGGER.info(f"Successfully processed {len(valid_results)} users")
    
    # 这里可以添加将结果写入S3的逻辑,比如拼接成CSV后上传
    csv_content = "FirstName,LastName,Username,RoutingProfile,SecurityProfiles,PhoneType,DeskPhoneNumber,AutoAccept,ACWTimeout,HierarchyGroup\r\n"
    csv_content += ''.join(valid_results)
    
    # 上传到S3
    s3_client = boto3.client('s3')
    try:
        s3_client.put_object(
            Bucket=s3_bucket,
            Key=f"connect_users_{datetime.now().strftime('%Y%m%d_%H%M%S')}.csv",
            Body=csv_content.encode('utf-8')
        )
        LOGGER.info(f"Successfully uploaded CSV to S3 bucket {s3_bucket}")
    except botocore.exceptions.ClientError as e:
        LOGGER.error(f"Failed to upload CSV to S3: {e}")
        return {'statusCode': 500, 'body': json.dumps('Failed to process users and upload to S3')}
    
    return {'statusCode': 200, 'body': json.dumps(f"Processed {len(valid_results)} users successfully")}

关键注意事项

  • 进程池大小:Lambda的vCPU数量决定了最优进程数,1vCPU的话开2-4个进程最合适,太多会导致资源竞争反而变慢
  • Boto3客户端:每个进程必须创建自己的客户端,因为boto3的客户端不是进程安全的,共享客户端会导致奇怪的错误
  • 错误处理:添加了详细的错误日志,避免单个用户处理失败导致整个任务崩溃
  • Lambda环境限制:Lambda的执行时间上限是15分钟,6000条用户如果每条处理需要1秒,并行4个进程的话大概需要25分钟,这会超时!如果是这种情况,你需要考虑分批处理(比如用SQS队列,每次处理1000条),或者优化API调用(比如检查是否有批量获取用户详情的API,不过Amazon Connect好像没有)
  • 模块级代码:Lambda冷启动时会执行模块级代码,所以不要在模块级创建进程池,要在lambda_handler里用with语句创建,确保每次调用都有干净的进程池

这样应该就能解决你的并行化问题,同时修复原代码里的错误啦!

内容的提问来源于stack exchange,提问作者user2373210

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.27 14:47:32