如何在Tweepy的callStream中捕获TweetListener的on_error返回的420错误
解决Twitter Streaming API 420错误时切换API密钥的问题
我之前也遇到过类似的420错误(Twitter的速率限制/临时封禁提示),要让上层的callStream函数感知到TweetListener里的错误,核心是建立错误回调机制,同时维护一组API密钥池来切换使用。下面是具体的实现方案:
第一步:修改TweetListener类,添加错误回调
我们需要让Listener在捕获到420错误时,主动通知上层逻辑。可以通过给Listener传入一个回调函数来实现:
import tweepy import json class TweetListener(tweepy.StreamListener): def __init__(self, error_callback=None, userid=None, projectid=None): super().__init__() self.error_callback = error_callback # 保存错误回调函数 self.userid = userid self.projectid = projectid def on_connect(self): print("You are now connected to the streaming API.") def on_error(self, status_code): print(f'An Error has occured: {repr(status_code)}') # 仅在捕获到420错误时触发回调 if status_code == 420 and self.error_callback: self.error_callback(status_code) # 返回False会断开当前流连接,为重新连接做准备 return False def on_data(self, data): json_data = json.loads(data) print(json_data) # 注意:这里要返回True,否则流会自动停止 return True
第二步:修改callStream函数,实现密钥切换与重连逻辑
我们需要维护一组API密钥,在错误触发时切换到下一组,并用指数退避的方式等待后重新启动流:
import tweepy from APIs.StreamKafkaApi1 import TweetListener import time # 准备多组有效的API密钥,按顺序循环使用 API_KEY_POOL = [ { "consumer_key": "FIRST_CONSUMER_KEY", "consumer_secret": "FIRST_CONSUMER_SECRET", "access_token": "FIRST_ACCESS_TOKEN", "access_secret": "FIRST_ACCESS_SECRET" }, { "consumer_key": "SECOND_CONSUMER_KEY", "consumer_secret": "SECOND_CONSUMER_SECRET", "access_token": "SECOND_ACCESS_TOKEN", "access_secret": "SECOND_ACCESS_SECRET" }, # 可根据需求添加更多密钥组 ] hashtags = ["#ipl"] current_key_idx = 0 active_stream = None def handle_stream_error(status_code): global current_key_idx, active_stream print(f"触发{status_code}错误,准备切换API密钥...") # 切换到下一组密钥(循环使用) current_key_idx = (current_key_idx + 1) % len(API_KEY_POOL) # 关闭当前活跃的流连接 if active_stream: active_stream.disconnect() # 指数退避等待:420错误后不要立刻重连,避免加剧限制 wait_seconds = 2 ** current_key_idx print(f"等待{wait_seconds}秒后重新连接...") time.sleep(wait_seconds) # 重新启动流 callStream() def callStream(): global active_stream # 获取当前要使用的API密钥 current_keys = API_KEY_POOL[current_key_idx] print(f"正在使用第{current_key_idx+1}组API密钥") # 初始化认证与API实例 auth = tweepy.OAuthHandler(current_keys["consumer_key"], current_keys["consumer_secret"]) auth.set_access_token(current_keys["access_token"], current_keys["access_secret"]) api = tweepy.API(auth, wait_on_rate_limit=True) # 传入错误回调初始化Listener tweet_listener = TweetListener( error_callback=handle_stream_error, userid=userid, projectid=projectid ) # 创建并启动流 active_stream = tweepy.Stream(api.auth, tweet_listener) try: active_stream.filter(track=hashtags, async=True) except Exception as e: print(f"流启动失败:{str(e)}") handle_stream_error(500) if __name__ == "__main__": # 替换为你的实际userid和projectid userid = "your_user_id" projectid = "your_project_id" callStream()
关键逻辑说明
- 错误回调传递:通过给
TweetListener传入error_callback,让Listener能在捕获到420错误时通知上层逻辑执行切换操作。 - 密钥池循环:用
current_key_idx跟踪当前使用的密钥,切换时循环遍历密钥池,避免单组密钥耗尽配额。 - 指数退避等待:420错误是Twitter的临时限制,指数增长的等待时间能有效降低再次触发限制的概率。
- 流状态管理:用全局变量
active_stream跟踪当前活跃的流,确保切换时能正确关闭旧连接。
内容的提问来源于stack exchange,提问作者AshrithGande
相关产品推荐
相关产品推荐

