Flink批处理-简单案例-01


一、简单案例

<?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.api.java.DataSet;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.operators.DataSource;
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 FLink批处理测试案例
 * @menu
 */
public class FlinkBatch01App {

    public static void main(String[] args) throws Exception {
        //构建执行任务环境以及任务的启动的入口, 存储全局相关的参数
        ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
        //设置并行度,方便看到效果
        env.setParallelism(1);
        //相同类型元素的数据流 source
        //FLink1.12之后流批一体,不再使用DataSet了
        DataSet stringDS = env.fromElements("java,SpringBoot", "spring cloud,redis",
                                                    "kafka,课堂");
        stringDS.print("处理前");
        DataSet 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");
    }
}