这篇文章给大家分享的是有关flink batch dataset的示例代码的内容。小编觉得挺实用的,因此分享给大家做个参考,一起跟随小编过来看看吧。
package hgs.flink_lesson
import org.apache.flink.api.java.utils.ParameterTool
import org.apache.flink.streaming.api.scala._
import org.apache.flink.api.scala.ExecutionEnvironment
import org.apache.flink.core.fs.FileSystem.WriteMode
import org.apache.flink.api.common.accumulators.Accumulator
import org.apache.flink.api.common.accumulators.IntCounter
import scala.collection.immutable.List
import scala.collection.mutable.ListBuffer
import scala.collection.immutable.HashMap
//import StreamExecutionEnvironment.class
object WordCount {
def main(args: Array[String]): Unit = {
val params = ParameterTool.fromArgs(args)
//1.获得一个执行环境,如果是Streaming则换成StreamExecutionEnvironment
val env = ExecutionEnvironment.getExecutionEnvironment
//这样会得到当前环境下的配置
env.getConfig.setGlobalJobParameters(params)
println(params.get("input"))
println(params.get("output"))
val text = if(params.has("input")){
//2.加载或者创建初始化数据
env.readTextFile(params.get("input"))
}else{
println("Please specify the input file directory.")
return
}
println("lines "+text.count())
val ac = new IntCounter
//3.在数据上指明操作类型
val counts = text.flatMap{ _.toLowerCase().split("\\W+").filter{_.nonEmpty}}
//这里与spark的算子的groupBy有点不同,这边要用数组类似的下标来确定根据什么进行分组
.map{(_,1)}.groupBy(0).reduceGroup(it=>{
val tuple = it.next()
var cnt = tuple._2
val ch = tuple._1
while(it.hasNext){
cnt= cnt+it.next()._2
}
(ch,cnt)})
//指明计算后的数据结果放到哪个位置
//4.counts.print()
counts.writeAsCsv("file:/d:/re.txt", "\n", " ",WriteMode.OVERWRITE)
//5.触发程序执行
env.execute("Scala WordCount Example")
//
}
}
感谢各位的阅读!关于“flink batch dataset的示例代码”这篇文章就分享到这里了,希望以上内容可以对大家有一定的帮助,让大家可以学到更多知识,如果觉得文章不错,可以把它分享出去让更多的人看到吧!
亿速云「云服务器」,即开即用、新一代英特尔至强铂金CPU、三副本存储NVMe SSD云盘,价格低至29元/月。点击查看>>
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。
原文链接:http://blog.itpub.net/31506529/viewspace-2564530/