Flink流处理-简单案例-01
一、pom文件
<?xml version="1.0" encoding="UTF-8"?>
4.0.0
com.robots
robots-flink
1.0-SNAPSHOT
UTF-8
UTF-8
1.8
1.8
1.8
2.12
1.13.1
org.projectlombok
lombok
1.18.16
org.apache.flink
flink-clients_${scala.version}
${flink.version}
org.apache.flink
flink-scala_${scala.version}
${flink.version}
org.apache.flink
flink-java
${flink.version}
org.apache.flink
flink-streaming-scala_${scala.version}
${flink.version}
org.apache.flink
flink-streaming-java_${scala.version}
${flink.version}
org.slf4j
slf4j-log4j12
1.7.7
runtime
log4j
log4j
1.2.17
runtime
com.alibaba
fastjson
1.2.44
二、简单流处理代码
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 stringDS = env.fromElements("java,SpringBoot", "spring cloud,redis",
"kafka,课堂");
stringDS.print("处理前");
DataStream flatMapDS = stringDS.flatMap(new FlatMapFunction() {
@Override
public void flatMap(String value, Collector collector) throws Exception {
String [] arr = value.split(",");
for(String str : arr){
collector.collect(str);
}
}
});
//输出 sink
flatMapDS.print("处理后");
//DataStream需要调用execute,可以取个名称
env.execute("flat map job");
}
}