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

如何断开Tweepy流并动态更新流过滤规则?

问题描述

需要每15分钟更新一次Tweepy流的过滤规则,尝试先断开流再重启,但自定义的disconnect方法(仅设置self.running=False)无法终止流。后续尝试逐个添加规则并启动流,发现调用stream.filter()后后续代码完全不执行,无法完成所有规则的添加,相关代码如下:

import telebot
import tweepy
import time
bot = telebot.TeleBot()
api_key = 
api_key_secret = 
bearer_token = 
access_token = 
access_token_secret = 
client = tweepy.Client(bearer_token, api_key, api_key_secret, access_token, access_token_secret)
auth = tweepy.OAuth1UserHandler(api_key, api_key_secret, access_token, access_token_secret)
api = tweepy.API(auth)


class MyStream(tweepy.StreamingClient):
    def on_connect(self):
        print('Connected')

    def on_response(self, response):
        print(response)
        

stream = MyStream(bearer_token=bearer_token, wait_on_rate_limit=True)

rules = ['book', 'tree', 'word']

#create the stream
for rule in rules:
    stream.add_rules(tweepy.StreamRule(rule))
    print('Showing the rule')
    print(stream.get_rules().data)
    stream.filter(tweet_fields=["referenced_tweets"])
    # this part of the code no longer works.
    print('sleep 10 sec')
    time.sleep(10)
# this part not working too
print('Final Streaming Rules:')
print(stream.get_rules().data)

问题原因

  1. stream.filter()是阻塞式方法:调用该方法后,程序会一直停留在流处理的循环中,不会继续执行后续代码,所以循环里第一次调用stream.filter()后,后面的打印、休眠以及剩余规则的添加逻辑都无法执行。
  2. 自定义disconnect方法无效:Tweepy的StreamingClient运行状态由内部多线程和连接逻辑共同控制,仅修改self.running无法正确终止底层网络连接和处理线程,必须使用官方提供的终止方法。

解决方法

1. 批量添加规则后启动流(单次启动场景)

如果只需一次性添加所有规则再启动流,正确逻辑是先添加完所有规则,再调用一次stream.filter():

import telebot
import tweepy
import time
bot = telebot.TeleBot()
api_key = ""
api_key_secret = ""
bearer_token = ""
access_token = ""
access_token_secret = ""
client = tweepy.Client(bearer_token, api_key, api_key_secret, access_token, access_token_secret)
auth = tweepy.OAuth1UserHandler(api_key, api_key_secret, access_token, access_token_secret)
api = tweepy.API(auth)


class MyStream(tweepy.StreamingClient):
    def on_connect(self):
        print('Connected')

    def on_response(self, response):
        print(response)
        

stream = MyStream(bearer_token=bearer_token, wait_on_rate_limit=True)

rules = ['book', 'tree', 'word']

# 先批量添加所有规则
for rule in rules:
    stream.add_rules(tweepy.StreamRule(rule))
    print('Added rule:', rule)

# 查看最终规则
print('Final Streaming Rules:')
print(stream.get_rules().data)

# 启动流
stream.filter(tweet_fields=["referenced_tweets"])

2. 定时更新规则(15分钟更新一次场景)

要实现定时更新规则,需要用线程启动流避免阻塞主程序,然后定时调用stream.disconnect()终止流,更新规则后重新启动:

import telebot
import tweepy
import time
import threading

bot = telebot.TeleBot()
api_key = ""
api_key_secret = ""
bearer_token = ""
access_token = ""
access_token_secret = ""
client = tweepy.Client(bearer_token, api_key, api_key_secret, access_token, access_token_secret)
auth = tweepy.OAuth1UserHandler(api_key, api_key_secret, access_token, access_token_secret)
api = tweepy.API(auth)

class MyStream(tweepy.StreamingClient):
    def on_connect(self):
        print('Connected')

    def on_response(self, response):
        print(response)

def start_stream(stream):
    """启动流的线程函数"""
    stream.filter(tweet_fields=["referenced_tweets"])

stream = MyStream(bearer_token=bearer_token, wait_on_rate_limit=True)
# 初始规则
current_rules = ['book', 'tree', 'word']

# 初始化规则
for rule in current_rules:
    stream.add_rules(tweepy.StreamRule(rule))

# 启动流线程
stream_thread = threading.Thread(target=start_stream, args=(stream,), daemon=True)
stream_thread.start()

# 定时更新逻辑(每15分钟)
update_interval = 15 * 60  # 15分钟秒数
while True:
    time.sleep(update_interval)
    
    # 停止流
    stream.disconnect()
    print('Stream disconnected, updating rules...')
    
    # 这里可根据需求修改规则,示例替换为新规则列表
    new_rules = ['new_book', 'new_tree', 'new_word']
    
    # 删除旧规则
    existing_rules = stream.get_rules().data
    if existing_rules:
        rule_ids = [rule.id for rule in existing_rules]
        stream.delete_rules(rule_ids)
    
    # 添加新规则
    for rule in new_rules:
        stream.add_rules(tweepy.StreamRule(rule))
    
    # 重启流线程
    stream_thread = threading.Thread(target=start_stream, args=(stream,), daemon=True)
    stream_thread.start()
    print('Stream restarted with new rules')

关键注意点

  • 终止流必须使用stream.disconnect(),这是Tweepy官方提供的方法,能正确关闭连接和终止内部线程。
  • 启动流时用线程包裹,避免阻塞主程序的定时更新逻辑。
  • 更新规则前要先删除旧规则,避免规则叠加导致重复匹配。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 14:15:40