文章23 | 阅读 12028 | 点赞0
package com.gosuncn
import org.apache.flink.api.scala._
object WordCountBatchJob {
def main(args: Array[String]) {
val env = ExecutionEnvironment.getExecutionEnvironment
val lines: DataSet[String] = env.readTextFile("C:\\Users\\root\\Desktop\\data.txt")
lines.flatMap(_.split(" ")).map((_, 1)).groupBy(0).sum(1).writeAsText("C:\\Users\\root\\Desktop\\out").setParallelism(2)
env.execute("WordCountBatchJob")
}
}
内容来源于网络,如有侵权,请联系作者删除!