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

如何让Telegram Python Bot在run_polling时同步获取用户新推文

解决方案

1. 改用Tweepy异步流避免阻塞事件循环

python-telegram-bot 20.x完全基于asyncio实现,同步代码会直接阻塞事件循环,导致/help这类命令无法被处理。Tweepy 4.x提供了AsyncStream类,天生适配异步环境,不会和机器人的命令处理抢资源。

2. 同一事件循环并行运行两个任务

用asyncio.gather把机器人的run_polling和Tweepy流的启动任务绑定,让它们在同一个asyncio事件循环里并行执行,两边的逻辑互不干扰。

3. 异步安全的共享数据管理

用asyncio.Lock保护用户的订阅列表(比如用户ID到Twitter ID的映射字典),确保多任务读写时不会出现数据错乱。如果要分文件部署,把共享数据和锁放到单独的Python模块里即可——Python模块是单例模式,跨文件导入时会共享同一实例。

单文件完整代码示例

import asyncio
from telegram import Update
from telegram.ext import ApplicationBuilder, CommandHandler, ContextTypes
import tweepy

# 共享订阅数据:键是Telegram聊天ID,值是订阅的Twitter ID列表
user_subscriptions = {}
# 异步锁,保护订阅数据的读写安全
subs_lock = asyncio.Lock()

# 自定义Tweepy异步流处理类
class MyAsyncStream(tweepy.AsyncStream):
    async def on_tweet(self, tweet):
        # 过滤掉回复和转发,只处理原创推文
        if tweet.in_reply_to_status_id is not None or tweet.referenced_tweets:
            return
        
        # 找出所有订阅该推文作者的Telegram用户
        target_chats = []
        async with subs_lock:
            for chat_id, twitter_ids in user_subscriptions.items():
                if tweet.author_id in twitter_ids:
                    target_chats.append(chat_id)
        
        # 给每个目标用户发送推文内容
        for chat_id in target_chats:
            try:
                await application.bot.send_message(
                    chat_id=chat_id,
                    text=f"@{tweet.author.username} 发新推文了:\n{tweet.text}"
                )
            except Exception as e:
                print(f"给{chat_id}发消息失败:{e}")

# Telegram命令处理器
async def start(update: Update, context: ContextTypes.DEFAULT_TYPE):
    await update.message.reply_text("欢迎使用!用/add <Twitter ID>添加订阅,/help看所有命令。")

async def add(update: Update, context: ContextTypes.DEFAULT_TYPE):
    if not context.args:
        await update.message.reply_text("格式错了!应该是:/add 123456(替换成目标Twitter ID)")
        return
    
    try:
        twitter_id = int(context.args[0])
        chat_id = update.effective_chat.id
        
        async with subs_lock:
            if chat_id not in user_subscriptions:
                user_subscriptions[chat_id] = []
            if twitter_id not in user_subscriptions[chat_id]:
                user_subscriptions[chat_id].append(twitter_id)
                await update.message.reply_text(f"已订阅Twitter ID:{twitter_id}")
            else:
                await update.message.reply_text("这个ID已经在你的订阅列表里了")
    except ValueError:
        await update.message.reply_text("Twitter ID必须是纯数字!")

async def help_cmd(update: Update, context: ContextTypes.DEFAULT_TYPE):
    await update.message.reply_text("可用命令:\n/start - 初始化使用\n/add <Twitter ID> - 添加订阅\n/list - 查看已订阅ID\n/remove <Twitter ID> - 取消订阅")

async def main():
    global application
    # 初始化Telegram机器人
    application = ApplicationBuilder().token("你的Telegram机器人Token").build()
    
    # 注册命令处理器
    application.add_handler(CommandHandler("start", start))
    application.add_handler(CommandHandler("add", add))
    application.add_handler(CommandHandler("help", help_cmd))
    
    # 初始化Tweepy异步流
    tweepy_auth = tweepy.OAuth1UserHandler(
        "你的Twitter API Key",
        "你的Twitter API Secret",
        "你的Twitter Access Token",
        "你的Twitter Access Token Secret"
    )
    stream = MyAsyncStream(tweepy_auth)
    
    # 同时启动机器人和Tweepy流
    await asyncio.gather(
        application.run_polling(poll_interval=1),
        stream.filter(follow=[], tweet_fields=["author_id", "referenced_tweets"])
    )

if __name__ == "__main__":
    asyncio.run(main())

动态更新Tweepy流的关注列表

用户会随时添加新订阅,需要定期更新Stream的follow参数来监听新的账号。可以加个定时任务:

# 新增定时更新流关注列表的函数
async def update_stream_follow(stream):
    while True:
        await asyncio.sleep(60)  # 每分钟更新一次
        async with subs_lock:
            # 去重后获取所有已订阅的Twitter ID
            all_twitter_ids = list({tid for tids in user_subscriptions.values() for tid in tids})
        if all_twitter_ids:
            await stream.disconnect()
            await stream.filter(follow=all_twitter_ids, tweet_fields=["author_id", "referenced_tweets"])

# 修改main里的gather,加入这个定时任务
await asyncio.gather(
    application.run_polling(poll_interval=1),
    update_stream_follow(stream),
    stream.filter(follow=[], tweet_fields=["author_id", "referenced_tweets"])
)

分文件部署的同步方案

如果要把机器人和流逻辑拆分到不同文件,只需把共享数据单独抽出来:

  1. 创建shared_data.py:
import asyncio

user_subscriptions = {}
subs_lock = asyncio.Lock()
  1. Telegram机器人模块(telegram_bot.py)和Tweepy流模块(tweepy_stream.py)都导入这个文件的变量,就能实现跨文件数据同步。最后在main.py里用asyncio.gather同时启动两个模块的运行函数即可。

关键注意事项

  • 绝对不能用Tweepy的同步Stream,会直接卡死asyncio事件循环,导致机器人命令完全无响应。
  • 所有读写订阅数据的操作必须用async with subs_lock:包裹,避免异步环境下的数据竞争。
  • 如果需要持久化订阅数据,推荐用aiosqlite异步操作SQLite,不要用同步数据库库(比如原生sqlite3),同样会阻塞事件循环。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 04:31:01