Python异步并行处理Google Workspace API调用问题排查
解决Google Workspace API异步并发调用问题
核心问题分析
你的脚本按顺序执行,本质原因是逐个await每个子OU的用户查询任务,而非将所有异步任务提交到事件循环中同时执行。asyncio的await会阻塞当前协程直到任务完成,顺序调用自然会变成串行执行。
关键修改步骤
1. 确保API调用是真正的异步实现
Google Workspace官方的google-api-python-client默认是同步库,直接包装成async函数依然是串行逻辑。必须使用异步HTTP客户端(如aiohttp)配合异步认证工具(如gcloud-aio-auth)来调用API,才能实现真正的异步请求。
2. 使用asyncio.gather()实现并发
将所有子OU的用户查询任务收集为任务列表,通过asyncio.gather()一次性提交到事件循环并发执行,替代逐个await的串行逻辑。
修改示例代码对比
错误的串行执行代码(你的初始版本可能类似)
import asyncio from googleapiclient.discovery import build async def get_sub_ous(parent_ou): # 同步API包装成async,实际仍为串行 service = build('admin', 'directory_v1') result = service.orgunits().list(customerId='my_customer', parentOrgUnitPath=parent_ou, type='children').execute() return result.get('organizationUnits', []) async def get_users_in_ou(ou_path): service = build('admin', 'directory_v1') users = [] page_token = None while True: result = service.users().list(customerId='my_customer', orgUnitPath=ou_path, pageToken=page_token).execute() users.extend(result.get('users', [])) page_token = result.get('nextPageToken') if not page_token: break return users async def main(): sub_ous = await get_sub_ous('/parent-ou') # 逐个await导致串行执行 all_users = [] for ou in sub_ous: users = await get_users_in_ou(ou['orgUnitPath']) all_users.extend(users) print(f"Total users: {len(all_users)}") asyncio.run(main())
正确的并发执行修改版
import asyncio import aiohttp from gcloud.aio.auth import Token # 异步调用Google OU列表API async def fetch_sub_ous(session, token, parent_ou): url = f"https://admin.googleapis.com/admin/directory/v1/customer/my_customer/orgunits?parentOrgUnitPath={parent_ou}&type=children" headers = {'Authorization': f'Bearer {token}'} async with session.get(url, headers=headers) as resp: result = await resp.json() return result.get('organizationUnits', []) # 异步调用Google用户列表API(支持分页) async def fetch_users_in_ou(session, token, ou_path): users = [] page_token = None while True: url = f"https://admin.googleapis.com/admin/directory/v1/users?customer=my_customer&orgUnitPath={ou_path}" if page_token: url += f"&pageToken={page_token}" headers = {'Authorization': f'Bearer {token}'} async with session.get(url, headers=headers) as resp: result = await resp.json() users.extend(result.get('users', [])) page_token = result.get('nextPageToken') if not page_token: break return users async def main(): # 获取异步认证Token token = Token('admin.googleapis.com', service_file='path/to/service-account.json') await token.fetch() async with aiohttp.ClientSession() as session: # 第一步:获取子OU列表 sub_ous = await fetch_sub_ous(session, token.value, '/parent-ou') # 第二步:创建所有用户查询任务,通过gather并发执行 user_tasks = [fetch_users_in_ou(session, token.value, ou['orgUnitPath']) for ou in sub_ous] all_users_lists = await asyncio.gather(*user_tasks) # 合并所有子OU的用户结果 all_users = [] for users in all_users_lists: all_users.extend(users) print(f"Total users: {len(all_users)}") asyncio.run(main())
验证并发的方法
在fetch_users_in_ou函数中添加时间戳打印,比如:
print(f"[{asyncio.get_event_loop().time()}] Start fetching users for OU: {ou_path}") # ... 请求代码 ... print(f"[{asyncio.get_event_loop().time()}] Finish fetching users for OU: {ou_path}")
执行日志会显示多个OU的查询请求几乎同时发起,而非按顺序等待上一个任务完成后才启动下一个。
注意事项
- 确保服务账号拥有Google Workspace的组织单元读取权限和用户读取权限。
- Google API有配额限制,可通过
asyncio.Semaphore限制并发任务数,避免触发配额告警:semaphore = asyncio.Semaphore(5) # 限制同时执行5个任务 async def fetch_users_in_ou(session, token, ou_path): async with semaphore: # 原有逻辑
内容的提问来源于stack exchange,提问作者AdeQ
相关产品推荐
相关产品推荐

