Flink流處理-簡單案例-01
阿新 • • 發佈:2022-03-10
一、pom檔案
<?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"> <modelVersion>4.0.0</modelVersion> <groupId>com.robots</groupId> <artifactId>robots-flink</artifactId> <version>1.0-SNAPSHOT</version> <properties> <encoding>UTF-8</encoding> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <maven.compiler.source>1.8</maven.compiler.source> <maven.compiler.target>1.8</maven.compiler.target> <java.version>1.8</java.version> <scala.version>2.12</scala.version> <flink.version>1.13.1</flink.version> </properties> <dependencies> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <version>1.18.16</version> </dependency> <!--flink客戶端--> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients_${scala.version}</artifactId> <version>${flink.version}</version> </dependency> <!--scala版本--> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-scala_${scala.version}</artifactId> <version>${flink.version}</version> </dependency> <!--java版本--> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-java</artifactId> <version>${flink.version}</version> </dependency> <!--streaming的scala版本--> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-scala_${scala.version}</artifactId> <version>${flink.version}</version> </dependency> <!--streaming的java版本--> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_${scala.version}</artifactId> <version>${flink.version}</version> </dependency> <!--日誌輸出--> <dependency> <groupId>org.slf4j</groupId> <artifactId>slf4j-log4j12</artifactId> <version>1.7.7</version> <scope>runtime</scope> </dependency> <dependency> <groupId>log4j</groupId> <artifactId>log4j</artifactId> <version>1.2.17</version> <scope>runtime</scope> </dependency> <!--json依賴包--> <dependency> <groupId>com.alibaba</groupId> <artifactId>fastjson</artifactId> <version>1.2.44</version> </dependency> </dependencies> </project>
二、簡單流處理程式碼
import org.apache.flink.api.common.functions.FlatMapFunction; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.datastream.DataStreamSource; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.util.Collector; /** * @datetime 2022-03-09 上午9:47 * @desc * @menu */ public class Flink01App { public static void main(String[] args) throws Exception { //構建執行任務環境以及任務的啟動的入口, 儲存全域性相關的引數 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); //設定並行度 env.setParallelism(1); //相同型別元素的資料流 source DataStreamSource<String> stringDS = env.fromElements("java,SpringBoot", "spring cloud,redis", "kafka,課堂"); stringDS.print("處理前"); DataStream<String> flatMapDS = stringDS.flatMap(new FlatMapFunction<String, String>() { @Override public void flatMap(String value, Collector<String> collector) throws Exception { String [] arr = value.split(","); for(String str : arr){ collector.collect(str); } } }); //輸出 sink flatMapDS.print("處理後"); //DataStream需要呼叫execute,可以取個名稱 env.execute("flat map job"); } }