```scala
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")
}
}
```
- Flink簡介
- flink搭建standalone模式與測試
- flink提交任務(界面方式)
- Flink項目初始化
- Java版WordCount(匿名類)
- Java版WordCount(lambda)
- Scala版WordCount
- Java版WordCount[批處理]
- Scala版WordCount[批處理]
- 流處理非并行的Source
- 流處理可并行的Source
- kafka的Source
- Flink算子(Map,FlatMap,Filter)
- Flink算子KeyBy
- Flink算子Reduce和Max與Min
- addSink自定義Sink
- startNewChain和disableChaining
- 資源槽slotSharingGroup
- 計數窗口
- 滾動窗口
- 滑動窗口
- Session窗口
- 按照EventTime作為標準