python-telegram-bot开发中如何将阻塞函数改造为异步形式
你的判断完全正确:基于asyncio的python-telegram-bot(v20及以上版本)运行在单线程事件循环上,若在协程handler中直接调用同步阻塞的pandas计算逻辑,会卡住整个事件循环,所有后续用户请求必须等当前计算完成才能被处理,并发场景下会出现明显的排队、响应超时问题。
注意:直接给pandas函数加async def前缀并不能解决阻塞问题。pandas本身是CPU密集型的同步库,没有原生异步实现,这种改法只是给同步逻辑套了一层协程外壳,执行时依然会阻塞事件循环。
方案1:用事件循环自带的执行器跑同步pandas逻辑(改造成本最低,首选)
原理是把阻塞的pandas计算扔到事件循环外的独立工作线程/进程中执行,事件循环本身不被阻塞,可以正常响应其他用户请求,你原有写好的pandas处理逻辑完全不需要改动。
针对普通耗时(单次计算1s以内)的场景,用asyncio.to_thread扔到默认线程池即可,改造示例:import asyncio import pandas as pd from telegram import Update from telegram.ext import ApplicationBuilder, CommandHandler, ContextTypes # 原有同步pandas函数不需要做任何修改 def build_metro_schedule(gtfs_raw_data, line_id: str) -> pd.DataFrame: # 所有pandas处理逻辑:GTFS解析、表关联、时间筛选、时刻表生成 # 保持原有代码不动即可 ... return result_df async def schedule_handler(update: Update, context: ContextTypes.DEFAULT_TYPE): line_id = context.args[0] # 把同步计算扔到独立线程执行,await过程不会阻塞事件循环 schedule_df = await asyncio.to_thread(build_metro_schedule, gtfs_data, line_id) reply_text = convert_df_to_text(schedule_df) await update.message.reply_text(reply_text)如果你的pandas计算量极大(单次计算超过1s、处理数据量超过十万行),线程池会受Python GIL限制性能不足,可以换成全局进程池绕开GIL:
import asyncio from concurrent.futures import ProcessPoolExecutor # 初始化全局进程池,避免每次请求重复创建进程,worker数和CPU核心数持平即可 process_pool = ProcessPoolExecutor(max_workers=4) async def heavy_schedule_handler(update: Update, context: ContextTypes.DEFAULT_TYPE): loop = asyncio.get_running_loop() line_id = context.args[0] # 注意:传入进程池的参数必须是可序列化的纯数据,不能传Update/Context这类不可序列化的对象 schedule_df = await loop.run_in_executor( process_pool, build_metro_schedule, gtfs_data, line_id ) await update.message.reply_text(convert_df_to_text(schedule_df))方案2:增加缓存层从根源减少重复计算(收益最高)
地铁GTFS时刻表、运行状态都属于短时间内不会频繁变动的数据,完全没必要每次用户请求都从头跑pandas计算:- 固定的基础运营时刻表,在机器人启动时一次性预计算完成存在内存中,用户请求直接读取结果
- 实时到站、运行状态类的计算结果,按线路、行驶方向做带过期时间的缓存,一般设置10-30秒有效期即可,大量重复请求直接返回缓存结果,不用重复执行计算
简单的内存缓存可以直接用cachetools的TTLCache实现,不需要额外部署缓存服务:
from cachetools import TTLCache from datetime import timedelta # 最多缓存1000条结果,每条有效期30秒 schedule_cache = TTLCache(maxsize=1000, ttl=timedelta(seconds=30).total_seconds()) def build_metro_schedule(gtfs_raw_data, line_id: str, direction: int) -> pd.DataFrame: cache_key = f"{line_id}_{direction}" if cache_key in schedule_cache: return schedule_cache[cache_key] # 原有pandas计算逻辑 ... result = ... schedule_cache[cache_key] = result return result方案3:替换为支持异步的数据分析库(不推荐,成本极高)
可以用Dask、Polars等带异步接口的数据分析库替换pandas,但你现有重度依赖pandas的业务代码需要全量重写,投入产出比极低,只有单次要处理GB级以上GTFS数据的场景才需要考虑。
- 不要在传给线程池/进程池的同步pandas函数里调用任何协程方法,会抛出事件循环状态错误
- 不要把线程池/进程池的worker数设得过大,线程池大小建议设为CPU核心数*2,进程池大小和CPU核心数持平即可,过多worker会因为频繁上下文切换拖慢整体性能
- 用进程池时不要跨进程传递pandas DataFrame对象,会大幅增加序列化开销,尽量在进程内完成数据计算后只返回最终需要的文本、结构化小结果
内容的提问来源于stack exchange,提问作者f3d3

