如何让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"]) )
分文件部署的同步方案
如果要把机器人和流逻辑拆分到不同文件,只需把共享数据单独抽出来:
- 创建
shared_data.py:
import asyncio user_subscriptions = {} subs_lock = asyncio.Lock()
- 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
相关产品推荐
相关产品推荐

