如何为基于Twython的Twitter流采集程序添加定时停止功能
定时停止功能实现方案
完全可以通过datetime模块实现该需求,无需引入额外依赖,仅需要修改自定义流类和初始化逻辑即可,具体修改如下:
修改步骤
- 给
MyStreamer类新增初始化方法,记录程序启动时间和最大允许运行时长 - 在
on_success回调方法最前端添加时长校验逻辑,运行超时则主动断开流连接 - 主函数初始化流实例时传入你需要的运行时长(单位为秒,如需运行1小时则传3600)
修改后的完整代码
from twython import Twython, TwythonError, TwythonStreamer import pandas as pd from datetime import datetime import csv import os APP_KEY = '***' APP_SECRET = '***' OAUTH_TOKEN = '***' OAUTH_TOKEN_SECRET = '***' # OAuth 2 twitter = Twython(APP_KEY, APP_SECRET, oauth_version=2) ACCESS_TOKEN = twitter.obtain_access_token() twitter = Twython(APP_KEY, access_token=ACCESS_TOKEN) # OAuth 1 twitter2 = Twython(APP_KEY, APP_SECRET) auth = twitter2.get_authentication_tokens(callback_url='https://twitter.com') OAUTH_TOKEN2 = auth['oauth_token'] OAUTH_TOKEN_SECRET2 = auth['oauth_token_secret'] def process_tweet(tweet): filtered_data = {} dict_test = [] initial_format = '%a %b %d %H:%M:%S %z %Y' date_formats = '%d-%m-%Y' time_formats = '%H:%M:%S' filtered_data['id_post'] = tweet['id'] filtered_data['hashtags'] = [hashtag['text'] for hashtag in tweet['entities']['hashtags']] filtered_data['date'] = datetime.strptime(tweet['created_at'], initial_format).strftime(date_formats) filtered_data['time'] = datetime.strptime(tweet['created_at'], initial_format).strftime(time_formats) filtered_data['geo'] = tweet['geo'] filtered_data['text'] = tweet['text'] filtered_data['user'] = tweet['user']['screen_name'] filtered_data['user_loc'] = tweet['user']['location'] filtered_data['user_id'] = tweet['in_reply_to_user_id'] filtered_data['source_device'] = tweet['source'] dict_test.append(filtered_data) print(dict_test) return dict_test class MyStreamer(TwythonStreamer): # 新增:初始化方法接收最大运行时长参数,记录启动时间 def __init__(self, app_key, app_secret, oauth_token, oauth_token_secret, max_run_seconds): super().__init__(app_key, app_secret, oauth_token, oauth_token_secret) self.start_time = datetime.now() self.max_run_seconds = max_run_seconds def on_success(self, data): # 新增:超时校验逻辑 run_duration = (datetime.now() - self.start_time).total_seconds() if run_duration > self.max_run_seconds: print("已达到设定运行时长,停止采集") self.disconnect() return # 原有业务逻辑不变 if data['lang'] == 'ru': tweet_data = process_tweet(data) self.save_to_csv(tweet_data) def on_error(self, status_code, data): # 适配Python3写法,Python2可保留原print语句 print(status_code) def save_to_csv(self, tweet, encoding = 'utf-8'): file_name = 'Twitter_{date}.csv'.format(date = str(datetime.now().strftime('%d-%m-%Y'))) try: with open(file_name, 'a', newline='', encoding=encoding) as file: size_path_file = os.path.getsize(file_name) print(size_path_file) if size_path_file == 0: writer = csv.DictWriter(file, fieldnames=tweet[0].keys()) writer.writeheader() for data in tweet: writer.writerow(data) else: writer = csv.DictWriter(file, fieldnames=tweet[0].keys()) for data in tweet: writer.writerow(data) except IOError: print("I/O error") if __name__ == '__main__': # 设定采集运行时长,单位为秒,此处3600为1小时,可自行修改 RUN_DURATION = 3600 stream = MyStreamer(APP_KEY, APP_SECRET, OAUTH_TOKEN, OAUTH_TOKEN_SECRET, RUN_DURATION) stream.statuses.filter(track=['Короновирус', 'Covid-19', 'Вакцинация', 'Вакцина', 'Карантин', 'Заболевание', 'Ограничения', 'Тесты', 'Ковид'])
Tweepy方案兼容性说明
核心的定时判断逻辑完全可以通用,仅需要注意两个库的流断开API差异:
- Tweepy v1版本的流对象断开方法为
close(),v2版本为disconnect() - Twython流的断开方法固定为
disconnect()
仅需调整对应断开调用的方法名,Tweepy的定时器逻辑就可以直接迁移到Twython中使用。
补充说明:如果你的采集场景可能长时间收不到推文导致超时判断不及时,可以额外开启一个守护线程做超时校验,普通流采集场景下上述方案的性能损耗几乎为0,完全够用。
内容的提问来源于stack exchange,提问作者kostya ivanov
相关产品推荐
相关产品推荐

