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

SpringBoot集成flink全過程

 更新時(shí)間:2025年01月14日 14:45:36   作者:A塵埃  
文章介紹了Flink作為批處理和流處理結(jié)合的統(tǒng)一計(jì)算框架,特別是其強(qiáng)大的流處理能力,文章還詳細(xì)描述了如何在本地和集群環(huán)境中搭建Flink,并通過Netcat工具生成一個(gè)無界流測試,文章最后提供了啟動socket流的步驟和示例代碼,希望對讀者有所幫助

SpringBoot集成flink

Flink是一個(gè)批處理和流處理結(jié)合的統(tǒng)一計(jì)算框架,其核心是一個(gè)提供了數(shù)據(jù)分發(fā)以及并行化計(jì)算的流數(shù)據(jù)處理引擎。

最大亮點(diǎn)是流處理,最適合的應(yīng)用場景是低時(shí)延的數(shù)據(jù)處理。

場景

高并發(fā)pipeline處理數(shù)據(jù),時(shí)延毫秒級,且兼具可靠性。

環(huán)境搭建

①、安裝flink

https://nightlies.apache.org/flink/flink-docs-master/zh/docs/try-flink/local_installation/

②、安裝Netcat

Netcat(又稱為NC)是一個(gè)計(jì)算機(jī)網(wǎng)絡(luò)工具,它可以在兩臺計(jì)算機(jī)之間建立 TCP/IP 或 UDP 連接。

用于測試網(wǎng)絡(luò)中的端口,發(fā)送文件等操作。

進(jìn)行網(wǎng)絡(luò)調(diào)試和探測,也可以進(jìn)行加密連接和遠(yuǎn)程管理等高級網(wǎng)絡(luò)操作

yum install -y nc # 安裝nc命令
nc -lk 8888 # 啟動socket端口

無界流之讀取socket文本流

一、依賴

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <parent>
        <artifactId>springboot-demo</artifactId>
        <groupId>com.et</groupId>
        <version>1.0-SNAPSHOT</version>
    </parent>
    <modelVersion>4.0.0</modelVersion>

    <artifactId>flink</artifactId>

    <properties>
        <maven.compiler.source>8</maven.compiler.source>
        <maven.compiler.target>8</maven.compiler.target>
    </properties>
    <dependencies>

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-autoconfigure</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-test</artifactId>
            <scope>test</scope>
        </dependency>
        <!-- 添加 Flink 依賴 -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-java</artifactId>
            <version>1.17.0</version>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-java -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-java</artifactId>
            <version>1.17.0</version>
        </dependency>

        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-clients -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-clients</artifactId>
            <version>1.17.0</version>
        </dependency>

        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-connector-base -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-connector-base</artifactId>
            <version>1.17.0</version>
        </dependency>

        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-connector-files</artifactId>
            <version>1.17.0</version>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-connector-kafka -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-connector-kafka</artifactId>
            <version>1.17.0</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-runtime-web</artifactId>
            <version>1.17.0</version>
        </dependency>


    </dependencies>
    <build>
        <plugins>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-shade-plugin</artifactId>
                <executions>
                    <execution>
                        <phase>package</phase>
                        <goals>
                            <goal>shade</goal>
                        </goals>
                        <configuration>
                            <transformers>
                                <transformer
                                        implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">
                                    <resource>META-INF/spring.handlers</resource>
                                </transformer>
                                <transformer
                                        implementation="org.springframework.boot.maven.PropertiesMergingResourceTransformer">
                                    <resource>META-INF/spring.factories</resource>
                                </transformer>
                                <transformer
                                        implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">
                                    <resource>META-INF/spring.schemas</resource>
                                </transformer>
                                <transformer
                                        implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer" />
                                <transformer
                                        implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
                                    <mainClass>com.et.flink.job.SocketJob</mainClass>
                                </transformer>
                            </transformers>
                        </configuration>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>

</project>

二、SoketJob

public class SocketJob{
	
	public static void main(String[] args)throws Exception{
		
		// 創(chuàng)建執(zhí)行環(huán)境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 指定并行度,默認(rèn)電腦線程數(shù)
        env.setParallelism(3);
        // 讀取數(shù)據(jù)socket文本流 指定監(jiān)聽 IP 端口 只有在接收到數(shù)據(jù)才會執(zhí)行任務(wù)
        DataStreamSource<String> socketDS = env.socketTextStream("172.24.4.193", 8888);

        // 處理數(shù)據(jù): 切換、轉(zhuǎn)換、分組、聚合 得到統(tǒng)計(jì)結(jié)果
        SingleOutputStreamOperator<Tuple2<String, Integer>> sum = socketDS
                .flatMap(
                        (String value, Collector<Tuple2<String, Integer>> out) -> {
                            String[] words = value.split(" ");
                            for (String word : words) {
                                out.collect(Tuple2.of(word, 1));
                            }
                        }
                )
                .setParallelism(2)
                // // 顯式地提供類型信息:對于flatMap傳入Lambda表達(dá)式,系統(tǒng)只能推斷出返回的是Tuple2類型,而無法得到Tuple2<String, Long>。只有顯式設(shè)置系統(tǒng)當(dāng)前返回類型,才能正確解析出完整數(shù)據(jù)
                .returns(new TypeHint<Tuple2<String, Integer>>() {
                })
//                .returns(Types.TUPLE(Types.STRING,Types.INT))
                .keyBy(value -> value.f0)
                .sum(1);


        // 輸出
        sum.print();

        // 執(zhí)行
        env.execute();
	}
}

測試:

啟動socket流:

nc -l 8888

本地執(zhí)行:直接ideal啟動main程序,在socket流中輸入

abc bcd cde
bcd cde fgh
cde fgh hij

集群執(zhí)行:

執(zhí)行maven打包,將打包的jar上傳到集群中

總結(jié)

以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • Spring?Initializr只能創(chuàng)建為Java?17版本以上的問題解決

    Spring?Initializr只能創(chuàng)建為Java?17版本以上的問題解決

    這篇文章主要給大家介紹了關(guān)于Spring?Initializr只能創(chuàng)建為Java?17版本以上問題的解決辦法,文中通過圖文介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2024-01-01
  • spring三級緩存以及為什么不用二級緩存解讀

    spring三級緩存以及為什么不用二級緩存解讀

    Spring三級緩存機(jī)制解決了循環(huán)依賴問題,通過一級緩存存放完全初始化的bean,二級緩存存放實(shí)例化但未完成依賴注入和初始化的bean,三級緩存存放bean的創(chuàng)建工廠,避免了重復(fù)創(chuàng)建和確保代理對象的正確生成
    2025-02-02
  • Tomcat調(diào)優(yōu)詳解

    Tomcat調(diào)優(yōu)詳解

    這篇文章主要介紹了Tomcat調(diào)優(yōu)方式,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2024-04-04
  • JAVA異常和自定義異常處理方式

    JAVA異常和自定義異常處理方式

    這篇文章主要介紹了JAVA異常和自定義異常處理方式,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-09-09
  • Springboot中的Shiro基礎(chǔ)入門教程

    Springboot中的Shiro基礎(chǔ)入門教程

    Apache Shiro是一個(gè)強(qiáng)大且靈活的開源安全框架,適用于Java應(yīng)用程序,它提供了認(rèn)證、授權(quán)、會話管理和加密等功能,通過Shiro,開發(fā)者可以輕松地實(shí)現(xiàn)用戶身份驗(yàn)證和權(quán)限控制,從而保護(hù)應(yīng)用程序的安全,本文介紹Springboot中的Shiro基礎(chǔ)入門教程,感興趣的朋友跟隨小編一起看看吧
    2025-12-12
  • Servlet實(shí)現(xiàn)簡單文件上傳功能

    Servlet實(shí)現(xiàn)簡單文件上傳功能

    這篇文章主要為大家詳細(xì)介紹了Servlet實(shí)現(xiàn)簡單文件上傳功能,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2017-10-10
  • 淺談web服務(wù)器項(xiàng)目中request請求和response的相關(guān)響應(yīng)處理

    淺談web服務(wù)器項(xiàng)目中request請求和response的相關(guān)響應(yīng)處理

    這篇文章主要介紹了淺談web服務(wù)器項(xiàng)目中request請求和response的相關(guān)響應(yīng)處理,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-07-07
  • Java SpringBoot自定義注解的使用及說明

    Java SpringBoot自定義注解的使用及說明

    本文介紹了在Spring Boot中創(chuàng)建和使用自定義注解的方法,通過自定義注解,可以減少重復(fù)代碼、增強(qiáng)代碼可讀性和可維護(hù)性,具體步驟包括定義注解、創(chuàng)建注解處理器以及在業(yè)務(wù)方法上使用注解
    2025-11-11
  • 別在Java代碼里亂打日志了,這才是正確的打日志姿勢

    別在Java代碼里亂打日志了,這才是正確的打日志姿勢

    這篇文章主要介紹了別在Java代碼里亂打日志了,這才是正確的打日志姿勢,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-06-06
  • java中Collection迭代器的實(shí)現(xiàn)

    java中Collection迭代器的實(shí)現(xiàn)

    Java迭代器是遍歷Collection集合的標(biāo)準(zhǔn)工具,提供hasNext()、next()和remove()三個(gè)核心方法,下面就來介紹一下Collection迭代器的實(shí)現(xiàn),感興趣的可以了解一下
    2026-01-01

最新評論

通化县| 富宁县| 上高县| 大冶市| 北碚区| 托克托县| 鹤岗市| 都兰县| 盐边县| 虎林市| 长武县| 胶南市| 措美县| 屯留县| 德安县| 色达县| 高邑县| 钟山县| 闵行区| 永泰县| 航空| 二连浩特市| 勃利县| 项城市| 西乌珠穆沁旗| 仁布县| 于都县| 成都市| 高淳县| 同江市| 讷河市| 福建省| 五常市| 晋城| 安阳县| 遂川县| 无为县| 林口县| 阿拉善盟| 朝阳市| 宝应县|