如何断开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)
问题原因
stream.filter()是阻塞式方法:调用该方法后,程序会一直停留在流处理的循环中,不会继续执行后续代码,所以循环里第一次调用stream.filter()后,后面的打印、休眠以及剩余规则的添加逻辑都无法执行。- 自定义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
相关产品推荐
相关产品推荐

