Python3.5+:如何将异步生成器合并为普通生成器?
嘿,这个问题我之前折腾过好一阵,异步生成器和普通生成器的衔接确实有点绕,不过咱们可以用「异步并行执行 + 普通生成器流式输出」的思路解决,既提速又不用把所有结果塞进内存。
第一步:把同步生成器改成异步生成器
首先得把你的google_search改成异步生成器(用async def定义,内部yield),这样才能用aiohttp异步发请求,实现并行爬取。示例代码大概是这样:
import aiohttp from bs4 import BeautifulSoup async def google_search_async(search_string): # 复用aiohttp会话,提升效率 async with aiohttp.ClientSession() as session: current_page = 1 while True: # 构造带分页参数的搜索URL(实际使用时建议用urllib.parse.quote处理搜索字符串编码) url = f"https://www.google.com/search?q={search_string}&start={(current_page-1)*10}" async with session.get(url, headers={"User-Agent": "Mozilla/5.0"}) as resp: if resp.status != 200: break # 请求失败就终止 html = await resp.text() soup = BeautifulSoup(html, "html.parser") results = soup.find_all("div", class_="g") if not results: break # 没有更多结果了 for result in results: # 提取你需要的信息(标题、链接等) title = result.find("h3").text if result.find("h3") else "无标题" link = result.find("a")["href"] if result.find("a") else "#" yield (title, link) # 异步yield结果 current_page += 1
第二步:用普通生成器包裹异步逻辑
现在要写一个普通生成器google_searches,内部利用asyncio事件循环来并行处理多个异步生成器,同时流式返回结果。这里分两种场景:
场景1:保持原zip的行为(返回各搜索的对应结果元组)
如果你需要和原来的zip逻辑一致——每次返回所有搜索的第n个结果组成的元组,那可以这样写:
import asyncio def google_searches(*search_strings): # 初始化所有异步生成器 async_gens = [google_search_async(q) for q in search_strings] loop = asyncio.get_event_loop() async def fetch_next(gen): try: # Python 3.6-3.9用gen.__anext__(),3.10+可以用anext(gen) return await gen.__anext__() except StopAsyncIteration: return None # 标记该生成器已结束 while True: # 为每个异步生成器创建「获取下一个结果」的任务 tasks = [loop.create_task(fetch_next(gen)) for gen in async_gens] # 等待所有任务完成(并行执行,比串行快N倍) results = loop.run_until_complete(asyncio.gather(*tasks)) # 只要有一个生成器结束,就终止循环(和zip逻辑一致) if None in results: break yield tuple(results)
使用的时候和原来一样:
for results in google_searches("python异步", "aiohttp教程"): # results是元组:(第一个搜索的结果, 第二个搜索的结果) print(results)
场景2:流式返回所有结果(哪个先爬完就先返回)
如果不需要严格对应各搜索的结果顺序,只想尽快拿到所有结果,那可以用asyncio.wait等待第一个完成的任务,实现真正的流式输出:
def google_searches_stream(*search_strings): async_gens = [google_search_async(q) for q in search_strings] loop = asyncio.get_event_loop() # 初始化任务字典:任务 → 对应的异步生成器 tasks = {loop.create_task(gen.__anext__()): gen for gen in async_gens} while tasks: # 等待第一个完成的任务 done, _ = loop.run_until_complete(asyncio.wait(tasks.keys(), return_when=asyncio.FIRST_COMPLETED)) for task in done: gen = tasks.pop(task) try: result = task.result() yield result # 直接返回这个结果 # 为该生成器创建下一个结果的任务 next_task = loop.create_task(gen.__anext__()) tasks[next_task] = gen except StopAsyncIteration: # 该搜索已无结果,跳过 pass
使用时直接遍历即可:
for result in google_searches_stream("python异步", "aiohttp教程"): print(result) # 哪个搜索先拿到结果就先打印
关键说明
- 为什么这样可行?普通生成器内部可以反复调用
loop.run_until_complete,每次等待一组异步任务完成,然后返回结果,既利用了异步并行的速度,又保持了普通生成器的流式特性,不用把所有结果存进内存。 - 资源清理:异步生成器里用
async with管理ClientSession,确保请求完成后关闭会话,避免资源泄漏。 - Python版本兼容:如果你用的是3.6,记得用
gen.__anext__()而不是anext(gen)(后者是3.10新增的)。
内容的提问来源于stack exchange,提问作者Max Smith
相关产品推荐
相关产品推荐

