异步生成器能否并发运行?及其实用场景解析
我想要更好地理解异步生成器(async generators)及其实用场景,因此编写了如下测试代码:
import asyncio async def download(urls): for url in urls: print(f"Downloading page at {url} started") await asyncio.sleep(2) # simulating a page download print(f"Page downloaded from {url}") yield f"Page downloaded from url {url}" async def main(): results = [] async for i in download(["foo.com", "bar.com", "baz.com"]): results.append(i) print("All results:", results) asyncio.run(main())
该代码的输出如下:
Downloading page at foo.com started Page downloaded from foo.com Downloading page at bar.com started Page downloaded from bar.com Downloading page at baz.com started Page downloaded from baz.com All results: ['Page downloaded from url foo.com', 'Page downloaded from url bar.com', 'Page downloaded from url baz.com']
可以看到页面并未被并发下载。我知道不使用异步生成器时,可通过以下代码创建独立协程实现并发下载:
async def download(url): print(f"Downloading page at {url} started") await asyncio.sleep(2) # simulating a page download print(f"Page downloaded from {url}") return f"Page downloaded from url {url}" async def main(): await asyncio.gather(*[download(x) for x in ["foo.com", "bar.com", "baz.com"]]) asyncio.run(main())
但我不清楚异步生成器的实用价值。能否结合代码示例解释异步生成器的实用场景?
异步生成器的核心价值不在于并发执行多个独立任务(这是asyncio.gather的强项),而在于异步地、逐个地产生结果,尤其适合处理需要流式处理、分批返回或迭代异步数据源的场景。以下是几个典型实用场景:
1. 流式处理异步数据源
当你需要从持续产生数据的异步源(比如WebSocket消息流、异步日志读取、分页API滚动查询)中逐个获取数据并即时处理时,异步生成器是绝佳选择。
比如模拟WebSocket消息流的实时处理场景:
import asyncio import random async def websocket_message_stream(): # 模拟WebSocket持续推送消息 for i in range(5): await asyncio.sleep(random.uniform(0.5, 2)) # 消息间隔随机 message = f"实时消息 {i+1}: 服务器状态更新" yield message async def process_stream(): async for msg in websocket_message_stream(): print(f"收到并处理消息: {msg}") asyncio.run(process_stream())
这里异步生成器持续产生消息,async for循环可以在每个消息产生后立即处理,无需等待所有数据生成完毕,完美适配实时业务场景。
2. 分批异步任务,逐个返回结果
如果你需要执行一批异步任务,但希望每完成一个任务就立即返回结果(而非等所有任务完成),同时还要控制并发数,可结合异步生成器与asyncio.Semaphore实现。
比如控制并发数的分批下载,每个结果就绪就返回:
import asyncio async def download_single(url, semaphore): async with semaphore: print(f"开始下载 {url}") await asyncio.sleep(2) # 模拟下载耗时 print(f"{url} 下载完成") return f"{url} 的结果" async def batch_download(urls, max_concurrent=2): semaphore = asyncio.Semaphore(max_concurrent) tasks = [download_single(url, semaphore) for url in urls] for task in asyncio.as_completed(tasks): result = await task yield result async def main(): async for result in batch_download(["url1", "url2", "url3", "url4"]): print(f"已获取结果: {result}") asyncio.run(main())
这个示例中,异步生成器batch_download通过asyncio.as_completed逐个获取完成的任务结果,既控制了并发数,又能实时反馈任务进度,适合需要即时展示结果的场景。
3. 迭代无限/大型异步数据集
如果要处理的数据集非常大甚至无限(比如异步读取大型日志文件、持续的传感器数据),一次性加载所有数据到内存不现实,异步生成器可以每次只生成一个数据项,大幅节省内存。
比如异步读取大型日志文件并筛选错误日志:
import asyncio async def read_large_log_file(file_path): # 模拟异步读取大型日志文件 with open(file_path, "r") as f: for line in f: await asyncio.sleep(0.01) # 模拟异步IO延迟 yield line.strip() async def process_logs(): async for log_line in read_large_log_file("large_log.txt"): if "ERROR" in log_line: print(f"发现错误日志: {log_line}") asyncio.run(process_logs())
这里异步生成器逐行读取日志,每读取一行就交给处理逻辑,无需把整个文件加载到内存,适合处理海量数据场景。
简单总结:异步生成器的优势是异步迭代+流式输出,当你需要边生成边处理数据、实时获取结果,或者处理无法一次性加载的异步数据源时,它比asyncio.gather更合适。
内容的提问来源于stack exchange,提问作者knightcool

