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

如何为基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 06:24:03