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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 13:23:10