【发布时间】:2016-06-26 16:57:59
【问题描述】:
我正在编写一个简单的字数统计 flink 作业,但我不断收到此错误:
could not find implicit value for evidence parameter of type org.apache.flink.api.common.typeinfo.TypeInformation[String]
[error] .flatMap{_.toLowerCase.split("\\W+") filter {_.nonEmpty}}
我搜索了网络,但没有得到任何可以理解的答案。
这是我的代码:
object Job {
def main(args: Array[String]) {
// set up the execution environment
val env = StreamExecutionEnvironment.getExecutionEnvironment
val dataStream = env.readTextFile("file:///home/plivo/code/flink/scala/flinkstream/test/")
val count = dataStream
.flatMap{_.toLowerCase.split("\\W+") filter {_.nonEmpty}}
.map{ (_,1) }
.groupBy(0)
.sum(1)
dataStream.print()
env.execute("Flink Scala API Skeleton")
}
}
【问题讨论】:
-
试试这个问题的答案,它也可能对你有帮助:stackoverflow.com/questions/29540121/…
-
我已经导入了所有必要的库,包括 flink.api.scala._ 和 flink.streaming.api.scala._
-
问题是flink(1.0.3版)的DataStream[(String, Int)]上没有groupBy(...)方法。有一个 keyBy(Int) 方法会产生一个 KeyedStream[(String, Int), Tuple]。
-
您能否尝试删除
import flink.api.scala._,因为流式传输以及批处理scala 包对象导入createTypeInformation。所以这些导入可能会发生冲突。
标签: scala apache-flink flink-streaming