SparkStreaming对接Kafka时reduceByKey方法调用报错求助
DStream.reduceByKey编译错误的问题 首先来看你遇到的编译错误:
KafkaCardCount.scala:56:28: value reduceByKey is not a member of org.apache.spark.streaming.dstream.DStream[Any]
[error] val wordCounts = etl.reduceByKey(_ + _)
[error] ^
[error] one error found
[error] (compile:compileIncremental) Compilation failed
[error] Total time: 7 s, completed Jan 14, 2018 2:52:23 PM
问题根因
你的etl DStream被Scala编译器推断为DStream[Any]类型了,这是因为在map操作里,你返回了两种完全不同的类型:
- 当
numStr符合数字格式时,返回(id, num)(属于Tuple2[String, Int]类型) - 当格式不匹配时,返回
()(属于Unit类型)
Scala会自动取这两种类型的公共父类Any作为整个DStream的元素类型。而reduceByKey是PairDStreamFunctions提供的方法,只有当DStream的元素是(Key, Value)这种键值对类型时才能调用,所以编译器会报错找不到这个方法。
修复方案
你需要过滤掉无效的数据,而不是返回Unit。这里推荐两种方式:
方式1:使用flatMap过滤无效元素
flatMap可以接受返回Option或集合的函数,自动过滤掉None或空集合的元素:
val etl = stream.flatMap(r => { val split = r.value.split("\t") val id = split(1) val numStr = split(4) if (numStr.matches("\\d+")) { val num = numStr.toInt Some((id, num)) // 用Some包装有效键值对 } else { None // 无效数据返回None,会被flatMap自动过滤 } })
方式2:map + filter组合处理
先提取数据,再过滤无效项,最后转换类型:
val etl = stream.map(r => { val split = r.value.split("\t") val id = split(1) val numStr = split(4) (id, numStr) }) // 过滤掉非数字的numStr .filter { case (_, numStr) => numStr.matches("\\d+") } // 将合法的numStr转为Int .map { case (id, numStr) => (id, numStr.toInt) }
修改后,etl的类型会被正确推断为DStream[(String, Int)],此时调用reduceByKey(_ + _)就可以正常编译了。
关于依赖配置
你的build.sbt配置是没问题的:
- Spark核心、Streaming以及Kafka 0-10的依赖版本匹配(都是2.1.0,对应Scala 2.11.8)
- 合并策略也正确处理了META-INF文件的冲突,避免打包时的错误
内容的提问来源于stack exchange,提问作者Yohei Onishi

