Flink是什么
Apache Flink 是一個框架和分布式處理引擎,用于對無界和有界資料流進行狀態計算,
Flink的特點
- 支持事件時間(event-time)和處理時間(processing-time)語意
- 精確一次(exactly-once)的狀態一致性保證
- 低延遲,每秒處理數百萬個事件,毫秒級延遲
- 與眾多常用存盤系統的連接
- 高可用,動態擴展,實作7*24小時全天候運行
Flink的全球熱度

Flink可以實作的目標
-
低延遲 來一次處理一次
-
高吞吐
-
結果的準確性和良好的容錯性
基于流的世界觀
- 在Flink的世界觀中,一切皆有流組成,就如python中的一切皆物件的概念,對應離線的資料,則規劃為有界流;對于實時的資料怎規劃為沒有界限的流,也就是Flink中的有界流于無界流
- 有開始也有結束的確定在一定時間范圍內的流稱為有界流,一旦確定就不會再改變,一般 批處理 用來處理有界資料,
- 無界流就是持續產生的資料流,資料是無限的,有開始,無結束,一般 流處理 用來處理無界資料

Flink第一課,三種方式實作詞頻統計
創建Flink工程
創建一個普通的maven工程,匯入相關依賴
<dependencies>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>1.10.1</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java_2.12</artifactId>
<version>1.10.1</version>
</dependency>
</dependencies>
匯入成功之后有一點要注意,就是java_2.12中的2.12指的是scala的版本,匯入依賴成功之后即在對應目錄創建包與對應類開始專案的撰寫,
批處理實作詞頻統計
package com.yo.wc;
/**
* created by YO
*/
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.DataSet;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.operators.DataSource;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.util.Collector;
// 批處理word count
public class WordCount {
public static void main(String[] args) throws Exception{
// 創建執行環境,類似與spark的創建背景關系
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
// 從檔案中讀取資料 這里可以隨意指定路徑,txt檔案寫入空格隔開的隨意單詞即可
String inputPath = "D:\\hello.txt";
//read讀取資料,可以指定讀取的檔案型別,整套批處理的api在flink里面就叫做dataset
//dataset是flink針對離線資料的處理模型
DataSet<String> inputDataSet = env.readTextFile(inputPath);
// 對資料集進行處理,按空格分詞展開,轉換成(word, 1)二元組進行統計
DataSet<Tuple2<String, Integer>> result = inputDataSet.flatMap(new MyFlatMapper())
.groupBy(0) // 按照第一個位置的word分組
.sum(1); // 將第二個位置上的資料求和
result.print();
}
// 自定義類,實作FlatMapFunction介面 輸出是String 輸出是元組Tuple2<String, Integer>>是flink提供的元組型別
public static class MyFlatMapper implements FlatMapFunction<String, Tuple2<String, Integer>> {
@Override
//value是輸入,out就是輸出的資料
public void flatMap(String value, Collector<Tuple2<String, Integer>> out) throws Exception {
// 按空格分詞
String[] words = value.split(" ");
// 遍歷所有word,包成二元組輸出
for (String word : words) {
out.collect(new Tuple2<>(word, 1));
}
}
}
}
輸出: 文本內的單詞不同輸出也不同
(scala,1)
(flink,1)
(world,1)
(hello,4)
流處理api實作詞頻統計
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.net.URL;
public class StreamWordCount {
public static void main(String[] args) throws Exception{
// 創建流處理執行環境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// // 從檔案中讀取資料
String inputPath = "D:\\hello.txt";
DataStream<String> inputDataStream = env.readTextFile(inputPath);
// 基于資料流進行轉換計算
DataStream<Tuple2<String, Integer>> resultStream = inputDataStream.flatMap(new WordCount.MyFlatMapper())
.keyBy(0)
.sum(1);
resultStream.print();
// 執行任務
env.execute();
}
}
輸出:

使用socket的方式
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.net.URL;
public class StreamWordCount {
public static void main(String[] args) throws Exception{
// 創建流處理執行環境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 用parameter tool工具從程式啟動引數中提取配置項 ,這里就是從main方法中獲取引數了args,可以在集群運行,這里再IDEA傳參模擬
ParameterTool parameterTool = ParameterTool.fromArgs(args);
String host = parameterTool.get("host");
int port = parameterTool.getInt("port");
// 從socket文本流讀取資料
DataStream<String> inputDataStream = env.socketTextStream(host, port);
// 基于資料流進行轉換計算
DataStream<Tuple2<String, Integer>> resultStream = inputDataStream.flatMap(new WordCount.MyFlatMapper())
.keyBy(0)
.sum(1);
resultStream.print();
// 執行任務
env.execute();
}
}

Flink的第一課入門到這里就完成了,同學們有遇到問題可直接私信,博主會盡力解答!
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/298318.html
標籤:其他
