在异步FastAPI应用的同步函数中调用异步函数的方法
问题描述
我有一个默认的异步FastAPI应用,其中包含一个同步函数:
@classmethod def make_values( cls, records: list, symbol: str = None, ) -> list: values = [] for record in records: new_row = SomePydanticModel( timestamp=record._mapping.get("time"), close=record._mapping.get("close"), value=record._mapping.get("value"), ) if symbol in EXTRA_SYMBOLS: new_row = split_row(new_row, symbol) values.append(new_row) return values
问题在于split_row是一个调用其他异步函数的异步函数:
async def split_row( row, symbol, ): adjusted_datetime = datetime.timestamp(datetime.strptime("01/01/19", "%m/%d/%y")) splits = await get_splits(symbol) # some big part with business logic with other async calls return row
目前new_row变量中得到的是协程对象,请问有没有办法获取split_row的执行结果?
解决方案
1. 将make_values改为异步函数(推荐)
这是最适配异步FastAPI场景的方案,直接在异步上下文里用await获取异步函数结果:
@classmethod async def make_values( cls, records: list, symbol: str = None, ) -> list: values = [] for record in records: new_row = SomePydanticModel( timestamp=record._mapping.get("time"), close=record._mapping.get("close"), value=record._mapping.get("value"), ) if symbol in EXTRA_SYMBOLS: new_row = await split_row(new_row, symbol) # 添加await获取结果 values.append(new_row) return values
之后调用make_values的地方也要用await,如果是FastAPI路由直接调用,把路由也改成异步即可(FastAPI原生支持异步路由)。
2. 在同步函数中手动运行事件循环
如果无法修改make_values的同步属性,可以获取当前运行的事件循环,强制执行异步函数:
import asyncio @classmethod def make_values( cls, records: list, symbol: str = None, ) -> list: values = [] loop = asyncio.get_running_loop() for record in records: new_row = SomePydanticModel( timestamp=record._mapping.get("time"), close=record._mapping.get("close"), value=record._mapping.get("value"), ) if symbol in EXTRA_SYMBOLS: # 运行异步任务并阻塞直到获取结果 new_row = loop.run_until_complete(split_row(new_row, symbol)) values.append(new_row) return values
注意:这种方式会阻塞事件循环,高并发场景下会拖慢FastAPI性能,仅适合小批量、低耗时的场景。
3. 批量异步处理提升效率
如果records数量较多,用asyncio.gather批量执行所有split_row任务,比逐个await更高效:
@classmethod async def make_values( cls, records: list, symbol: str = None, ) -> list: values = [] tasks = [] for record in records: new_row = SomePydanticModel( timestamp=record._mapping.get("time"), close=record._mapping.get("close"), value=record._mapping.get("value"), ) if symbol in EXTRA_SYMBOLS: tasks.append(split_row(new_row, symbol)) values.append(None) # 占位 else: values.append(new_row) # 批量执行异步任务并替换占位值 if tasks: results = await asyncio.gather(*tasks) idx = 0 for i in range(len(values)): if values[i] is None: values[i] = results[idx] idx +=1 return values
这种方式能充分利用异步IO的优势,减少整体等待时间。
内容的提问来源于stack exchange,提问作者Headmaster
相关产品推荐
相关产品推荐

