最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

淺談實(shí)時(shí)計(jì)算框架Flink集群搭建與運(yùn)行機(jī)制

 更新時(shí)間:2021年06月23日 15:47:50   作者:知了一笑  
Flink是一個(gè)框架和分布式處理引擎,用于對(duì)無(wú)界和有界數(shù)據(jù)流進(jìn)行有狀態(tài)計(jì)算。Flink被設(shè)計(jì)在所有常見(jiàn)的集群環(huán)境中運(yùn)行,以?xún)?nèi)存執(zhí)行速度和任意規(guī)模來(lái)執(zhí)行計(jì)算

一、Flink概述

1.1、基礎(chǔ)簡(jiǎn)介

主要特性包括:批流一體化、精密的狀態(tài)管理、事件時(shí)間支持以及精確一次的狀態(tài)一致性保障等。Flink不僅可以運(yùn)行在包括YARN、Mesos、Kubernetes在內(nèi)的多種資源管理框架上,還支持在裸機(jī)集群上獨(dú)立部署。在啟用高可用選項(xiàng)的情況下,它不存在單點(diǎn)失效問(wèn)題。

這里要說(shuō)明兩個(gè)概念:

  • 邊界:無(wú)邊界和有邊界數(shù)據(jù)流,可以理解為數(shù)據(jù)的聚合策略或者條件;
  • 狀態(tài):即執(zhí)行順序上是否存在依賴(lài)關(guān)系,即下次執(zhí)行是否依賴(lài)上次結(jié)果;

1.2、應(yīng)用場(chǎng)景

Data Driven

事件驅(qū)動(dòng)型應(yīng)用無(wú)須查詢(xún)遠(yuǎn)程數(shù)據(jù)庫(kù),本地?cái)?shù)據(jù)訪問(wèn)使得它具有更高的吞吐和更低的延遲,以反欺詐案例來(lái)看,DataDriven把處理的規(guī)則模型寫(xiě)到DatastreamAPI中,然后將整個(gè)邏輯抽象到Flink引擎,當(dāng)事件或者數(shù)據(jù)流入就會(huì)觸發(fā)相應(yīng)的規(guī)則模型,一旦觸發(fā)規(guī)則中的條件后,DataDriven會(huì)快速處理并對(duì)業(yè)務(wù)應(yīng)用進(jìn)行通知。

Data Analytics

和批量分析相比,由于流式分析省掉了周期性的數(shù)據(jù)導(dǎo)入和查詢(xún)過(guò)程,因此從事件中獲取指標(biāo)的延遲更低。不僅如此,批量查詢(xún)必須處理那些由定期導(dǎo)入和輸入有界性導(dǎo)致的人工數(shù)據(jù)邊界,而流式查詢(xún)則無(wú)須考慮該問(wèn)題,F(xiàn)link為持續(xù)流式分析和批量分析都提供了良好的支持,實(shí)時(shí)處理分析數(shù)據(jù),應(yīng)用較多的場(chǎng)景如實(shí)時(shí)大屏、實(shí)時(shí)報(bào)表。

Data Pipeline

與周期性的ETL作業(yè)任務(wù)相比,持續(xù)數(shù)據(jù)管道可以明顯降低將數(shù)據(jù)移動(dòng)到目的端的延遲,例如基于上游的StreamETL進(jìn)行實(shí)時(shí)清洗或擴(kuò)展數(shù)據(jù),可以在下游構(gòu)建實(shí)時(shí)數(shù)倉(cāng),確保數(shù)據(jù)查詢(xún)的時(shí)效性,形成高時(shí)效的數(shù)據(jù)查詢(xún)鏈路,這種場(chǎng)景在媒體流的推薦或者搜索引擎中十分常見(jiàn)。

二、環(huán)境部署

2.1、安裝包管理

[root@hop01 opt]# tar -zxvf flink-1.7.0-bin-hadoop27-scala_2.11.tgz

[root@hop02 opt]# mv flink-1.7.0 flink1.7

2.2、集群配置

管理節(jié)點(diǎn)

[root@hop01 opt]# cd /opt/flink1.7/conf

[root@hop01 conf]# vim flink-conf.yaml

jobmanager.rpc.address: hop01

分布節(jié)點(diǎn)

[root@hop01 conf]# vim slaves

hop02

hop03

兩個(gè)配置同步到所有集群節(jié)點(diǎn)下面。

2.3、啟動(dòng)與停止

/opt/flink1.7/bin/start-cluster.sh

/opt/flink1.7/bin/stop-cluster.sh

啟動(dòng)日志:

[root@hop01 conf]# /opt/flink1.7/bin/start-cluster.sh

Starting cluster.

Starting standalonesession daemon on host hop01.

Starting taskexecutor daemon on host hop02.

Starting taskexecutor daemon on host hop03.

2.4、Web界面

訪問(wèn):http://hop01:8081/

三、開(kāi)發(fā)入門(mén)案例

3.1、數(shù)據(jù)腳本

分發(fā)一個(gè)數(shù)據(jù)腳本到各個(gè)節(jié)點(diǎn):

/var/flink/test/word.txt

3.2、引入基礎(chǔ)依賴(lài)

這里基于Java寫(xiě)的基礎(chǔ)案例。

<dependencies>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-java</artifactId>
        <version>1.7.0</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-java_2.11</artifactId>
        <version>1.7.0</version>
    </dependency>
</dependencies>

3.3、讀取文件數(shù)據(jù)

這里直接讀取文件中的數(shù)據(jù),經(jīng)過(guò)程序流程分析出每個(gè)單詞出現(xiàn)的次數(shù)。

public class WordCount {
    public static void main(String[] args) throws Exception {
        // 讀取文件數(shù)據(jù)
        readFile () ;
    }

    public static void readFile () throws Exception {
        // 1、執(zhí)行環(huán)境創(chuàng)建
        ExecutionEnvironment environment = ExecutionEnvironment.getExecutionEnvironment();

        // 2、讀取數(shù)據(jù)文件
        String filePath = "/var/flink/test/word.txt" ;
        DataSet<String> inputFile = environment.readTextFile(filePath);

        // 3、分組并求和
        DataSet<Tuple2<String, Integer>> wordDataSet = inputFile.flatMap(new WordFlatMapFunction(
        )).groupBy(0).sum(1);

        // 4、打印處理結(jié)果
        wordDataSet.print();
    }

    // 數(shù)據(jù)讀取個(gè)切割方式
    static class WordFlatMapFunction implements FlatMapFunction<String, Tuple2<String, Integer>> {
        @Override
        public void flatMap(String input, Collector<Tuple2<String, Integer>> collector){
            String[] wordArr = input.split(",");
            for (String word : wordArr) {
                collector.collect(new Tuple2<>(word, 1));
            }
        }
    }
}

3.4、讀取端口數(shù)據(jù)

在hop01服務(wù)上創(chuàng)建一個(gè)端口,并模擬一些數(shù)據(jù)發(fā)送到該端口:

[root@hop01 ~]# nc -lk 5566

c++,java

通過(guò)Flink程序讀取并分析該端口的數(shù)據(jù)內(nèi)容:

public class WordCount {
    public static void main(String[] args) throws Exception {
        // 讀取端口數(shù)據(jù)
        readPort ();
    }

    public static void readPort () throws Exception {
        // 1、執(zhí)行環(huán)境創(chuàng)建
        StreamExecutionEnvironment environment = StreamExecutionEnvironment.getExecutionEnvironment();

        // 2、讀取Socket數(shù)據(jù)端口
        DataStreamSource<String> inputStream = environment.socketTextStream("hop01", 5566);

        // 3、數(shù)據(jù)讀取個(gè)切割方式
        SingleOutputStreamOperator<Tuple2<String, Integer>> resultDataStream = inputStream.flatMap(
                new FlatMapFunction<String, Tuple2<String, Integer>>()
        {
            @Override
            public void flatMap(String input, Collector<Tuple2<String, Integer>> collector) {
                String[] wordArr = input.split(",");
                for (String word : wordArr) {
                    collector.collect(new Tuple2<>(word, 1));
                }
            }
        }).keyBy(0).sum(1);

        // 4、打印分析結(jié)果
        resultDataStream.print();

        // 5、環(huán)境啟動(dòng)
        environment.execute();
    }
}

四、運(yùn)行機(jī)制

4.1、FlinkClient

客戶(hù)端用來(lái)準(zhǔn)備和發(fā)送數(shù)據(jù)流到JobManager節(jié)點(diǎn),之后根據(jù)具體需求,客戶(hù)端可以直接斷開(kāi)連接,或者維持連接狀態(tài)等待任務(wù)處理結(jié)果。

4.2、JobManager

在Flink集群中,會(huì)啟動(dòng)一個(gè)JobManger節(jié)點(diǎn)和至少一個(gè)TaskManager節(jié)點(diǎn),JobManager收到客戶(hù)端提交的任務(wù)后,JobManager會(huì)把任務(wù)協(xié)調(diào)下發(fā)到具體的TaskManager節(jié)點(diǎn)去執(zhí)行,TaskManager節(jié)點(diǎn)將心跳和處理信息發(fā)送給JobManager。

4.3、TaskManager

任務(wù)槽(slot)是TaskManager中最小的資源調(diào)度單位,在啟動(dòng)的時(shí)候就設(shè)置好了槽位數(shù),每個(gè)槽位能啟動(dòng)一個(gè)Task,接收J(rèn)obManager節(jié)點(diǎn)部署的任務(wù),并進(jìn)行具體的分析處理。

五、源代碼地址

GitHub·地址

https://github.com/cicadasmile/big-data-parent

GitEE·地址

https://gitee.com/cicadasmile/big-data-parent

以上就是淺談實(shí)時(shí)計(jì)算框架Flink集群搭建與運(yùn)行機(jī)制的詳細(xì)內(nèi)容,更多關(guān)于實(shí)時(shí)計(jì)算框架 Flink集群搭建與運(yùn)行機(jī)制的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • 安裝Ubuntu 16.04后要做的事(總結(jié))

    安裝Ubuntu 16.04后要做的事(總結(jié))

    Ubuntu 16.04發(fā)布了,帶來(lái)了很多新特性,同樣也依然帶著很多不習(xí)慣的東西,所以裝完系統(tǒng)后還要進(jìn)行一系列的優(yōu)化。本篇文章主要介紹了安裝Ubuntu 16.04后要做的事,有興趣的可以了解一下。
    2016-12-12
  • Linux之進(jìn)程狀態(tài)&&進(jìn)程優(yōu)先級(jí)詳解

    Linux之進(jìn)程狀態(tài)&&進(jìn)程優(yōu)先級(jí)詳解

    文章介紹了操作系統(tǒng)中進(jìn)程的狀態(tài),包括運(yùn)行狀態(tài)、阻塞狀態(tài)和掛起狀態(tài),并詳細(xì)解釋了Linux下進(jìn)程的具體狀態(tài)及其管理,此外,文章還討論了進(jìn)程的優(yōu)先級(jí)、查看和修改進(jìn)程優(yōu)先級(jí)的方法,以及并發(fā)相關(guān)的概念和函數(shù)的返回值
    2025-02-02
  • 在服務(wù)器上啟用HTTP公鑰固定擴(kuò)展的教程

    在服務(wù)器上啟用HTTP公鑰固定擴(kuò)展的教程

    這篇文章主要介紹了在服務(wù)器上啟用HTTP公鑰固定擴(kuò)展的教程,示例包括對(duì)Apache和NGINX以及Lighttpd服務(wù)器的演示,需要的朋友可以參考下
    2015-06-06
  • Linux初始化系統(tǒng)盤(pán)后重新掛載數(shù)據(jù)盤(pán)方法

    Linux初始化系統(tǒng)盤(pán)后重新掛載數(shù)據(jù)盤(pán)方法

    在本篇文章中我們給大家分享了Linux初始化系統(tǒng)盤(pán)后重新掛載數(shù)據(jù)盤(pán)的解決方法,有需要的朋友們可以參考下。
    2018-09-09
  • ssh修改超時(shí)自動(dòng)登出時(shí)間的方法

    ssh修改超時(shí)自動(dòng)登出時(shí)間的方法

    這篇文章主要介紹了關(guān)于linux中ssh超時(shí)自動(dòng)登出時(shí)間的設(shè)置方法,以避免總是被強(qiáng)行退出。需要的朋友可以參考借鑒,下面來(lái)一起看看吧。
    2017-02-02
  • Linux bzip2 命令的使用

    Linux bzip2 命令的使用

    這篇文章主要介紹了Linux bzip2 命令的使用,幫助大家更好的理解和使用Linux系統(tǒng),感興趣的朋友可以了解下
    2020-08-08
  • linux引導(dǎo)系統(tǒng)的方法分析

    linux引導(dǎo)系統(tǒng)的方法分析

    這篇文章主要介紹了linux引導(dǎo)系統(tǒng)的方法,總結(jié)分析了Linux引導(dǎo)系統(tǒng)相關(guān)原理、操作命令與注意事項(xiàng),需要的朋友可以參考下
    2020-03-03
  • Nginx 0.7.x + PHP 5.2.6(FastCGI)+ MySQL 5.1 在128M小內(nèi)存VPS服務(wù)器上的配置優(yōu)化

    Nginx 0.7.x + PHP 5.2.6(FastCGI)+ MySQL 5.1 在128M小內(nèi)存VPS服務(wù)器上的

    VPS(全稱(chēng)Virtual Private Server)是利用最新虛擬化技術(shù)在一臺(tái)物理服務(wù)器上創(chuàng)建多個(gè)相互隔離的虛擬私有主機(jī)。它們以最大化的效率共享硬件、軟件許可證以及管理資源。
    2008-12-12
  • RHEL 7中防火墻的配置和使用方法

    RHEL 7中防火墻的配置和使用方法

    下面小編就為大家?guī)?lái)一篇RHEL 7中防火墻的配置和使用方法。小編覺(jué)得挺不錯(cuò)的,現(xiàn)在就分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2016-12-12
  • 從Centos7升級(jí)到Centos8的教程(圖文詳解)

    從Centos7升級(jí)到Centos8的教程(圖文詳解)

    這篇文章主要介紹了從Centos7升級(jí)到Centos8的教程,在升級(jí)之前需要配置備份,本文通過(guò)圖文并茂的形式給大家介紹的非常詳細(xì),具有一定的參考借鑒價(jià)值,需要的朋友參考下吧
    2019-11-11

最新評(píng)論

巨野县| 贵港市| 城口县| 滦平县| 尉犁县| 封丘县| 嘉荫县| 清新县| 沾化县| 托克托县| 阜阳市| 黄大仙区| 麟游县| 余姚市| 十堰市| 三江| 镇赉县| 宜州市| 宿迁市| 连州市| 修文县| 敦煌市| 科技| 乐陵市| 会同县| 化州市| 金溪县| 得荣县| 海南省| 万安县| 云梦县| 石家庄市| 黎城县| 砚山县| 五台县| 湾仔区| 阿勒泰市| 安塞县| 余姚市| 临潭县| 琼海市|