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

Zeppelin中Spark Streaming调用TwitterUtils报Logging类缺失错误求助

解决Spark 2.x中TwitterUtils.createStream的NoClassDefFoundError: org/apache/spark/Logging问题

问题原因

你碰到的这个错误完全是版本不兼容导致的:你当前用的spark-streaming-twitter_2.10-1.0.0.jar是针对Spark 1.x + Scala 2.10开发的,它依赖的org.apache.spark.Logging类在Spark 2.0及以上版本已经被移除(Spark 2.x改用SLF4J作为日志接口),所以在Spark 2.2.1环境下运行必然会抛出类找不到的异常。另外这个jar包的Scala版本是2.10,和你当前用的Scala 2.11也不匹配,会带来额外的潜在问题。

解决方案

这里提供两种可行的解决方式,你可以根据自己的情况选择:

方式一:使用兼容Spark 2.x的Twitter Streaming组件

Spark官方在2.x之后不再维护spark-streaming-twitter,但Apache Bahir项目提供了兼容Spark 2.x的版本:

  • 移除旧jar包:在Zeppelin笔记本中删掉spark-streaming-twitter_2.10-1.0.0.jar
  • 添加兼容jar包:加入对应Scala 2.11和Spark 2.2.1版本的Bahir包:org.apache.bahir:spark-streaming-twitter_2.11:2.2.1
  • 验证代码:你原有的val tweets = TwitterUtils.createStream(ssc, None)代码不需要修改,因为Bahir的API和官方旧版本保持一致,导入的包还是org.apache.spark.streaming.twitter.TwitterUtils

方式二:自定义Twitter Receiver(更灵活,无第三方依赖)

如果不想依赖额外的Spark组件包,可以直接基于twitter4j实现自定义Receiver,完全避开版本兼容问题:

  1. 精简jar包:只保留twitter4j-core-4.0.4.jar和twitter4j-stream-4.0.4.jar,删掉spark-streaming-twitter_2.10-1.0.0.jar
  2. 实现自定义Receiver:在Zeppelin中添加以下代码:
import org.apache.spark.streaming.receiver.Receiver
import twitter4j._
import twitter4j.conf.ConfigurationBuilder
import org.apache.spark.storage.StorageLevel

class TwitterReceiver extends Receiver[Status](StorageLevel.MEMORY_AND_DISK_2) {
  override def onStart(): Unit = {
    // 读取你通过setUpTwitter设置的OAuth配置
    val cb = new ConfigurationBuilder()
    cb.setOAuthConsumerKey(System.getProperty("twitter4j.oauth.consumerKey"))
    cb.setOAuthConsumerSecret(System.getProperty("twitter4j.oauth.consumerSecret"))
    cb.setOAuthAccessToken(System.getProperty("twitter4j.oauth.accessToken"))
    cb.setOAuthAccessTokenSecret(System.getProperty("twitter4j.oauth.accessTokenSecret"))
    
    // 创建Twitter Stream并添加状态监听器
    val twitterStream = new TwitterStreamFactory(cb.build()).getInstance()
    val statusListener = new StatusListener {
      override def onStatus(status: Status): Unit = store(status)
      override def onDeletionNotice(notice: StatusDeletionNotice): Unit = {}
      override def onTrackLimitationNotice(limit: Int): Unit = {}
      override def onScrubGeo(userId: Long, upToId: Long): Unit = {}
      override def onStallWarning(warning: StallWarning): Unit = {}
      override def onException(ex: Exception): Unit = restart("重启Receiver", ex)
    }
    twitterStream.addListener(statusListener)
    // 如需过滤特定话题,可改用twitterStream.filter(new FilterQuery().track("你的关键词"))
    twitterStream.sample()
  }

  override def onStop(): Unit = {}
}
  1. 替换DStream创建代码:把原来的TwitterUtils.createStream调用换成:
val tweets = ssc.receiverStream(new TwitterReceiver())

额外注意事项

  • 确保你的twitter.txt文件中的OAuth配置项名称正确:应该是consumerKey、consumerSecret、accessToken、accessTokenSecret,和代码中System.getProperty的参数对应
  • Zeppelin中添加或删除jar包后,最好重启一下笔记本或者解释器,确保jar包加载生效

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:12:44