You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何我的异步代码性能未优于同步代码?求优化建议

异步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

问题排查与优化建议

核心问题点

  1. 阻塞性操作未异步化:async_build_and_assert_dataframe_from_csv函数虽标记为async,但内部的pd.read_csv是同步阻塞操作(CPU/IO密集),没有任何await调用,会占用事件循环线程,导致其他异步任务无法并行执行,完全浪费了异步框架的优势。
  2. 逻辑错误:空DataFrame判断逻辑完全反转——当前代码在len(df) > 0时打印"Dataframe is empty"警告,实际应该在len(df) == 0时触发。
  3. 冗余的任务存储:用set存储异步任务没有必要,直接用列表更简洁,不影响性能但增加代码复杂度。
  4. 响应读取方式冗余:手动调用result.decode("utf-8")可以用response.text()替代,aiohttp原生支持直接获取文本内容。

优化方案

  1. 将阻塞操作移至线程池:用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
    
  2. 简化任务创建逻辑:
    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)
    
  3. 优化API响应读取:
    async def api_get(self, session, url: str):
        async with session.get(url, headers=self.headers) as response:
            return await response.text()  # 直接获取文本,省去手动decode
    
  4. 并行处理DataFrame格式化:如果format_df操作可拆分,建议在process函数中完成单个子DataFrame的格式化,再进行最终合并,进一步利用并行性。

内容的提问来源于stack exchange,提问作者cluelessThrasher

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 19:51:16