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

使用Docker读取指定关键词Twitter推文时参数报错求助

问题定位与修复方案

问题背景

尝试通过Docker结合Kafka读取含指定关键词的Twitter推文,参考开源项目修改代码后,出现参数数量不匹配的类型错误,配置信息完整。

报错信息

参数数量错误截图
错误提示:TypeError: __init__() takes from 1 to 2 positional arguments but 3 were given

问题代码

### twitter
import tweepy
from tweepy.auth import OAuthHandler
from tweepy import Stream
import json
import logging 

### logging 
FORMAT = "%(asctime)s | %(name)s - %(levelname)s - %(message)s"
LOG_FILEPATH = "C:\\docker-kafka\\log\\testing.log"
logging.basicConfig(
    filename=LOG_FILEPATH,
    level=logging.INFO,
    filemode='w',
    format=FORMAT)

### Authenticate to Twitter
with open('C:\\docker-kafka\\credential.json','r') as f:
    credential = json.load(f)

CONSUMER_KEY = credential['twitter_api_key']
CONSUMER_SECRET = credential['twitter_api_secret_key']
ACCESS_TOKEN = credential['twitter_access_token']
ACCESS_TOKEN_SECRET = credential['twitter_access_token_secret']
BEARER_TOKEN = credential['bearer_token']

from kafka import KafkaProducer
producer = KafkaProducer(bootstrap_servers='localhost:9092',
    value_serializer=lambda v: v.encode('utf-8')) #Same port as your Kafka server

topic_name = "docker-twitter"

class twitterAuth():
    """SET UP TWITTER AUTHENTICATION"""
    def authenticateTwitterApp(self):
        auth = OAuthHandler(consumer_key=CONSUMER_KEY, consumer_secret=CONSUMER_SECRET)
        auth.set_access_token(ACCESS_TOKEN, ACCESS_TOKEN_SECRET)
        return auth

class TwitterStreamer():
    """SET UP STREAMER"""
    def __init__(self):
        self.twitterAuth = twitterAuth()

    def stream_tweets(self):
        while True:
            listener = ListenerTS() 
            auth = self.twitterAuth.authenticateTwitterApp()
            stream = Stream(auth, listener)
            stream.filter(track=["Starbucks"], stall_warnings=True, languages= ["en"])

class ListenerTS(tweepy.Stream):
    def on_status(self, status):
        tweet = json.dumps({
            'id': status.id, 
            'text': status.text, 
            'created_at': status.created_at.strftime("%Y-%m-%d %H:%M:%S")
        }, default=str)  
        producer.send(topic_name, tweet)
        return True

if __name__ == "__main__":
    TS = TwitterStreamer()
    TS.stream_tweets()

问题根源

错误出在Stream(auth, listener)这一行:

  • ListenerTS错误继承了tweepy.Stream类,而tweepy.Stream的构造函数不支持同时传入auth和listener两个参数(它的构造函数最多接受2个参数,这里实际传入了3个:self、auth、listener)。
  • 正确的监听类应该继承tweepy.StreamListener(对应Tweepy v1.1版本),作为回调类传递给tweepy.Stream实例。

修复步骤

  1. 修正监听类的继承关系
    将class ListenerTS(tweepy.Stream):改为:
class ListenerTS(tweepy.StreamListener):
  1. 添加错误处理方法
    为ListenerTS添加on_error方法,避免因API限制或其他错误导致程序中断:
def on_error(self, status_code):
    if status_code == 420:
        # 达到API请求限制,返回False断开连接
        return False
    logging.error(f"Twitter API Error: {status_code}")
    return True
  1. 清理冗余代码
    移除重复的from tweepy import OAuthHandler, Stream导入语句,以及未使用的BEARER_TOKEN变量。

修复后完整代码

import tweepy
from tweepy.auth import OAuthHandler
from tweepy import Stream
from tweepy.streaming import StreamListener
import json
import logging 
from kafka import KafkaProducer

### logging 
FORMAT = "%(asctime)s | %(name)s - %(levelname)s - %(message)s"
LOG_FILEPATH = "C:\\docker-kafka\\log\\testing.log"
logging.basicConfig(
    filename=LOG_FILEPATH,
    level=logging.INFO,
    filemode='w',
    format=FORMAT)

### Authenticate to Twitter
with open('C:\\docker-kafka\\credential.json','r') as f:
    credential = json.load(f)

CONSUMER_KEY = credential['twitter_api_key']
CONSUMER_SECRET = credential['twitter_api_secret_key']
ACCESS_TOKEN = credential['twitter_access_token']
ACCESS_TOKEN_SECRET = credential['twitter_access_token_secret']

producer = KafkaProducer(bootstrap_servers='localhost:9092',
    value_serializer=lambda v: v.encode('utf-8'))

topic_name = "docker-twitter"

class twitterAuth():
    """SET UP TWITTER AUTHENTICATION"""
    def authenticateTwitterApp(self):
        auth = OAuthHandler(consumer_key=CONSUMER_KEY, consumer_secret=CONSUMER_SECRET)
        auth.set_access_token(ACCESS_TOKEN, ACCESS_TOKEN_SECRET)
        return auth

class TwitterStreamer():
    """SET UP STREAMER"""
    def __init__(self):
        self.twitterAuth = twitterAuth()

    def stream_tweets(self):
        while True:
            listener = ListenerTS() 
            auth = self.twitterAuth.authenticateTwitterApp()
            stream = Stream(auth, listener)
            stream.filter(track=["Starbucks"], stall_warnings=True, languages= ["en"])

class ListenerTS(StreamListener):
    def on_status(self, status):
        tweet = json.dumps({
            'id': status.id, 
            'text': status.text, 
            'created_at': status.created_at.strftime("%Y-%m-%d %H:%M:%S")
        }, default=str)  
        producer.send(topic_name, tweet)
        logging.info(f"Sent tweet to Kafka: {status.id}")
        return True

    def on_error(self, status_code):
        if status_code == 420:
            return False
        logging.error(f"Twitter API Error Code: {status_code}")
        return True

if __name__ == "__main__":
    TS = TwitterStreamer()
    TS.stream_tweets()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 07:50:49