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

如何在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()

关键逻辑说明

  1. 错误回调传递:通过给TweetListener传入error_callback,让Listener能在捕获到420错误时通知上层逻辑执行切换操作。
  2. 密钥池循环:用current_key_idx跟踪当前使用的密钥,切换时循环遍历密钥池,避免单组密钥耗尽配额。
  3. 指数退避等待:420错误是Twitter的临时限制,指数增长的等待时间能有效降低再次触发限制的概率。
  4. 流状态管理:用全局变量active_stream跟踪当前活跃的流,确保切换时能正确关闭旧连接。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:15:52