基于Python Telegram Bot的多频道并发推送问题求助
问题背景
- 基于Python Telegram Bot(PTB)开发频道管理机器人,核心逻辑为抓取互联网帖子、处理后推送到多个目标频道
- 需求:为每个频道创建独立异步任务,每60秒刷新任务列表,新增频道自动创建对应任务
- 前期尝试用
threading实现出现事件循环冲突,改用AsyncIOScheduler+app.create_task方案后,频道推送正常,但操作机器人时频繁出现连接池耗尽错误
当前主调度代码
builder = ApplicationBuilder() builder.connection_pool_size(50000) builder.get_updates_connection_pool_size(50000) builder.pool_timeout(100) builder.get_updates_pool_timeout(100) app = builder.token(TOKEN).get_updates_http_version('1.1').http_version('1.1').build() thread_list = [] async def schedule(): channel_list = mydb.select("channels", {"visibility": 1, "active": "activate"}) for one in channel_list: if one['username'] not in thread_list: app.create_task(schedule_work(one['username'])) thread_list.append(one['username']) scheduler = AsyncIOScheduler() scheduler.add_job(schedule,'interval', seconds=60, name="main") scheduler.start() app.run_polling()
报错信息
.
.
File "/usr/lib/python3.10/asyncio/locks.py", line 214, in wait
await fut
asyncio.exceptions.CancelledError
.
.
File "/root/abaradmin/venv/lib/python3.10/site-packages/anyio/_core/_tasks.py", line 119, in exit
raise TimeoutError
TimeoutError
.
.
.
File "/root/abaradmin/venv/lib/python3.10/site-packages/httpcore/_exceptions.py", line 14, in map_exceptions
raise to_exc(exc) from exc
httpcore.PoolTimeout
.
.
File "/root/abaradmin/venv/lib/python3.10/site-packages/httpx/_transports/default.py", line 77, in map_httpcore_exceptions
raise mapped_exc(message) from exc
httpx.PoolTimeout
.
.
File "/root/abaradmin/venv/lib/python3.10/site-packages/telegram/request/_httpxrequest.py", line 226, in do_request
raise TimedOut(
telegram.error.TimedOut: Pool timeout: All connections in the connection pool are occupied. Request was not sent to Telegram. Consider adjusting the connection pool size or the pool timeout.
schedule_work函数代码
async def schedule_work(channel_usr): #time.sleep(60 * 5) sent_messages = {} start_time = datetime.datetime.now() while True: file_url = mydb.select('setting', {'section': 'channels_folder'})[1]['value'] file_dir = mydb.select('setting', {'section': 'channels_folder'})[0]['value'] sent_messages[channel_usr] = 0 # datas channel_username = channel_usr channel = mydb.select('channels', {"username": f"{channel_username}"}) schedule_channel = mydb.select('schedule', {"is_sent": 0, "channel": channel_username}) for post in schedule_channel: sending_type = channel[0]['sending_type'] sleep_status = channel[0]['sleep_status'] sleep_time = channel[0]['sleep_time'] channel_id = post['channel_id'] channel_admin = post['channel_admin'] post_type = post['type'] sleep_inquiry = is_sleep(sleep_time) if sleep_status != '✅ فعال' or sleep_inquiry == False: if post_type == 'music': # datas post_name = post['post_name'] mp3_file = str(post['mp3_file']) mp3_name = str(post['mp3_file']).split('/')[-1] mp3_file_128 = f'{file_dir}/{channel_username}/' + str(post['mp3_file']).split('/')[-1].split('.')[0] + ' [128].mp3' voice = str(post['voice']) cover_file = str(post['cover_file']) cover_name = str(post['cover_file']).split('/')[-1] caption = post['caption'] music_file_cap = "" if post['file_cap']: music_file_cap = post['file_cap'] voice_file_cap = "" if post['voice_cap']: voice_file_cap = post['voice_cap'] if sending_type == 'automatic': sent_messages[channel_username] += 3 # sending post try: await bot.send_photo(channel_id, f"{file_url}/{channel_username}/{cover_name}", caption=caption) time.sleep(0.3) await bot.send_voice(channel_id, open(voice, 'rb'), caption=voice_file_cap) time.sleep(0.3) await bot.send_audio(channel_id, f"{file_url}/{channel_username}/{mp3_name}", caption=music_file_cap) time.sleep(0.3) mydb.update('schedule', {'is_sent': 1, 'sent_date': datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")}, {'id': post["id"]}) # Deleting sent files os.system(f'rm -rf "{mp3_file}" "{mp3_file_128}"') os.system(f'rm -rf "{voice}" "{cover_file}"') except telegram.error.RetryAfter as e: log_writer(format_exc()) time.sleep(int(e.retry_after)) except telegram.error.TimedOut: log_writer(format_exc()) time.sleep(60) except: log_writer(format_exc()) await bot.send_message(log_channel, f'''Error in sending post at music automatic section: {format_exc()}''') else: if not post['admin_received']: sent_messages[channel_username] += 1 # sending post to admin for approval try: # Approval keyboard post_approval_keyboard = [ [ InlineKeyboardButton('🚫 حذف', callback_data=f'delete666_{post["id"]}'), InlineKeyboardButton('✅ انتشار', callback_data=f'share666_{post["id"]}'), ], [ InlineKeyboardButton('🚫 حذف همه', callback_data=f'delete_all_{channel_username}'), InlineKeyboardButton('✳️ انتشار همه', callback_data=f'share_all_{channel_username}'), ] ] post_approval = await bot.send_photo(channel_admin, f"{file_url}/{channel_username}/{cover_name}", caption=f"""تنظیمات نوع ارسال کانال @{channel_username} شما روی دستی تنظیم شده، موزیک «{post_name}» توی کانال منتشر بشه؟""", reply_markup=InlineKeyboardMarkup(post_approval_keyboard)) time.sleep(0.3) mydb.update('schedule', {'admin_received': 1, "is_confirmed": 0, "is_deleted": 0, 'approval_msg_id': post_approval['message_id']}, {'id': post["id"]}, True) except: log_writer(format_exc()) await bot.send_message(log_channel, f'''Error in sending post at music automatic section: {format_exc()}''') elif post_type == 'video': # bla bla bla time.sleep(0.7) end_time = datetime.datetime.now() finish_time = end_time - start_time print() print() print(channel_username, '>>>>>>>>>>', 'finish time:', finish_time, 'Sent message: ', sent_messages[channel_username] ) print() print() one_minute = datetime.timedelta(minutes=1) if int(sent_messages[channel_username]) >= 17: if finish_time < one_minute: print() print() print('********************* Sleep Mode *********************' ) print() print() time.sleep(40) sent_messages[channel_username] = 0 start_time = datetime.datetime.now() time.sleep(5)
解决方案
1. 替换同步sleep为异步sleep
代码中大量使用time.sleep()会阻塞整个事件循环,导致连接无法及时释放,直接耗尽连接池。所有同步sleep必须替换为异步版本:
import asyncio # 替换所有time.sleep(x)为await asyncio.sleep(x) await asyncio.sleep(0.3) await asyncio.sleep(60) await asyncio.sleep(40) await asyncio.sleep(5)
2. 限制并发任务数量
每个schedule_work是无限循环任务,频道数量增多时会导致并发任务过载,占用大量连接。通过信号量限制并发数:
# 在主代码中定义信号量,根据服务器性能和Telegram限制调整,建议设为10以内 semaphore = asyncio.Semaphore(8) async def schedule_work(channel_usr): async with semaphore: # 原函数逻辑不变,仅替换sleep为异步版本 ...
3. 修复任务生命周期管理
当前thread_list只添加不删除,频道停用再启用时无法重新创建任务,改用字典跟踪活跃任务并自动清理:
# 替换thread_list为任务字典,键为频道username,值为任务对象 active_tasks = {} async def schedule(): channel_list = mydb.select("channels", {"visibility": 1, "active": "activate"}) current_channels = {one['username'] for one in channel_list} # 清理已停用频道的任务 for username in list(active_tasks.keys()): if username not in current_channels: active_tasks[username].cancel() del active_tasks[username] # 为新增频道创建任务 for one in channel_list: username = one['username'] if username not in active_tasks: task = app.create_task(schedule_work(username)) active_tasks[username] = task
4. 优化数据库查询
schedule_work中每次循环重复查询相同配置,改为只查询一次:
async def schedule_work(channel_usr): async with semaphore: # 移到循环外,仅查询一次配置 setting_result = mydb.select('setting', {'section': 'channels_folder'}) file_url = setting_result[1]['value'] file_dir = setting_result[0]['value'] sent_messages = {} start_time = datetime.datetime.now() while True: sent_messages[channel_usr] = 0 # datas channel_username = channel_usr channel = mydb.select('channels', {"username": f"{channel_username}"}) schedule_channel = mydb.select('schedule', {"is_sent": 0, "channel": channel_username}) ...
5. 调整连接池配置
无需设置过大的连接池,异步环境下合理配置即可:
builder.connection_pool_size(50) builder.get_updates_connection_pool_size(10) builder.pool_timeout(30) builder.get_updates_pool_timeout(30)
6. 替换同步文件删除为异步操作
os.system('rm -rf ...')是同步阻塞操作,改用异步子进程:
# 替换同步删除 os.system(f'rm -rf "{mp3_file}" "{mp3_file_128}"') os.system(f'rm -rf "{voice}" "{cover_file}"') # 为异步版本 await asyncio.subprocess.create_subprocess_shell(f'rm -rf "{mp3_file}" "{mp3_file_128}"') await asyncio.subprocess.create_subprocess_shell(f'rm -rf "{voice}" "{cover_file}"')
内容的提问来源于stack exchange,提问作者Arta Alipour

