为何我的异步代码性能未优于同步代码?求优化建议
异步API数据管道性能优化排查
我搭建了一套数据管道,用于从API端点获取报告、执行基础清洗转换,并将数据合并为输出DataFrame。原本使用multiprocessing实现,但当前平台不再支持该方案,于是尝试async编程,对比了处理4份报告的n次运行中async、multithreading及同步循环调用的性能表现。但通过单因素方差分析(ANOVA 1-way)发现,异步代码在统计上并未比同步代码更快,怀疑是异步写法不够优化,恳请排查以下代码问题:
异步实现代码
class AsyncApiCaller(ApiCaller): # snip async def async_build_and_assert_dataframe_from_csv(self, data): df = pd.read_csv(StringIO(data)) if len(df) > 0: print("Warning: Dataframe is empty") return df async def process(self, session, url): result = await self.api_get(session, url) return await self.async_build_and_assert_dataframe_from_csv( result.decode("utf-8") ) async def make_all_requests(self, ids): async with aiohttp.ClientSession() as session: tasks = set() for url in ids: task = asyncio.create_task(coro=self.process(session=session, url=url)) tasks.add(task) return await asyncio.gather(*tasks, return_exceptions=False) async def api_get(self, session, url: str): async with session.get(url, headers=self.headers) as response: return await response.content.read() async def main(self) -> pd.DataFrame: # snip urls = # a list of urls to run results = await self.make_all_requests(urls) # concat results into single DataFrame final_df = pd.concat(results) final_df = self.format_df(final_df) # snip return final_df #snip
父类同步版本代码
class ApiCaller: # snip def go_single(self, id: ID) -> pd.DataFrame: url = self.get_url(id) # turn id into an endpoint url response = self.call_api(url, "GET", headers) # requests library wrapper data = self.parse_json_response(response, "data") dataframe = self.build_and_assert_dataframe_from_json(data) return dataframe def main(self) -> pd.DataFrame: # snip ids = # list of ids tasks = [self.go_single(id) for id in ids] dataframe = pd.concat(list(tasks), axis=0, ignore_index=False) dataframe = self.format_df(dataframe) # snip return dataframe
问题排查与优化建议
核心问题点
- 阻塞性操作未异步化:
async_build_and_assert_dataframe_from_csv函数虽标记为async,但内部的pd.read_csv是同步阻塞操作(CPU/IO密集),没有任何await调用,会占用事件循环线程,导致其他异步任务无法并行执行,完全浪费了异步框架的优势。 - 逻辑错误:空DataFrame判断逻辑完全反转——当前代码在
len(df) > 0时打印"Dataframe is empty"警告,实际应该在len(df) == 0时触发。 - 冗余的任务存储:用
set存储异步任务没有必要,直接用列表更简洁,不影响性能但增加代码复杂度。 - 响应读取方式冗余:手动调用
result.decode("utf-8")可以用response.text()替代,aiohttp原生支持直接获取文本内容。
优化方案
- 将阻塞操作移至线程池:用
asyncio.to_thread把pd.read_csv这类同步操作放到线程池执行,释放事件循环:async def async_build_and_assert_dataframe_from_csv(self, data): # 把阻塞的pd.read_csv放到线程池异步执行 df = await asyncio.to_thread(pd.read_csv, StringIO(data)) if len(df) == 0: print("Warning: Dataframe is empty") return df - 简化任务创建逻辑:
async def make_all_requests(self, ids): async with aiohttp.ClientSession() as session: tasks = [] for url in ids: task = asyncio.create_task(self.process(session=session, url=url)) tasks.append(task) return await asyncio.gather(*tasks, return_exceptions=False) - 优化API响应读取:
async def api_get(self, session, url: str): async with session.get(url, headers=self.headers) as response: return await response.text() # 直接获取文本,省去手动decode - 并行处理DataFrame格式化:如果
format_df操作可拆分,建议在process函数中完成单个子DataFrame的格式化,再进行最终合并,进一步利用并行性。
内容的提问来源于stack exchange,提问作者cluelessThrasher
相关产品推荐
相关产品推荐

