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,完全避开版本兼容问题:
- 精简jar包:只保留
twitter4j-core-4.0.4.jar和twitter4j-stream-4.0.4.jar,删掉spark-streaming-twitter_2.10-1.0.0.jar - 实现自定义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 = {} }
- 替换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
相关产品推荐
相关产品推荐

