Spark Scala访问Twitter推文流报错需改用V2流替代的问题咨询
报错产生原因
这个报错是Twitter官方API版本迭代导致的:你当前代码调用的是Twitter V1.1版本的采样流接口/1.1/statuses/sample.json,该接口目前已经正式下线,官方明确要求所有流类数据采集场景必须迁移到V2版本的过滤流、采样流接口,无论你的开发者账号是什么权限等级,调用已下线的V1.1流接口都会返回该提示。
绝大多数出现这个问题的场景,都是因为使用了停更多年的旧版依赖(比如老版本twitter4j、Spark Streaming原生的Twitter连接器),这类依赖默认的请求路径、鉴权逻辑都是适配V1.1版本API的,没有做V2适配,别费时间调试V1.1的参数、权限,这个接口是全量下线,任何等级的开发者账号都调不通。
解决方法
- 先完成开发者后台配置:进入Twitter开发者平台对应项目的设置页,确认App权限为只读及以上等级,在「Keys and Tokens」栏目生成V2版本专用的Bearer Token,V2流接口仅支持该凭证鉴权,V1.1版本使用的API Key/Secret、Access Token/Secret组合无法直接用于V2流接口。
- 替换采集依赖:移除项目中旧版twitter4j、Spark Streaming旧Twitter连接器这类未适配V2的依赖,可选择自行实现V2接口的HTTP长连接逻辑对接Spark,或使用已适配V2版本的Scala Twitter客户端。注意把请求路径替换为V2版本对应地址:
- 采样流地址:
https://api.twitter.com/2/tweets/sample/stream - 过滤流地址:
https://api.twitter.com/2/tweets/search/stream
- 采样流地址:
- 调整鉴权逻辑:V2流接口不需要走V1.1的OAuth1.0签名流程,发起HTTP请求时直接在请求头添加
Authorization: Bearer 你的V2 BearerToken即可。 - 调整请求参数:V2接口默认仅返回推文ID和文本内容,需要在请求URL后拼接参数指定你需要拉取的额外字段,常用参数格式参考:
?tweet.fields=created_at,lang,public_metrics&expansions=author_id&user.fields=username,verified,原V1.1接口的stall_warnings=true参数在V2接口中为默认开启状态,不需要额外传递。 - 适配Spark流处理逻辑:如果使用Spark Structured Streaming,简单测试可先将V2长连接拉取到的推文数据转发到本地Socket端口,再通过Socket Source接入,示例代码如下:
import org.apache.spark.sql.SparkSession val spark = SparkSession.builder() .appName("TwitterV2StreamCollect") .master("local[*]") .getOrCreate() // 接入本地转发的V2推文流 val tweetDF = spark.readStream .format("socket") .option("host", "localhost") .option("port", 9999) .load() // 编写流处理逻辑,比如分词、计数、写入存储等 val query = tweetDF.writeStream .outputMode("append") .format("console") .start() query.awaitTermination()
生产环境建议自行实现符合V2流规范的Structured Streaming自定义Source,内置断连重试、限流、背压处理逻辑即可。
内容的提问来源于stack exchange,提问作者Aaditya Mishra
相关产品推荐
相关产品推荐

