Flink的map/flatMap/filter的开发示例

Flink算子的开发,需要创建Maven项目,构建jar包,在Flink任务启动时,加载Jar包进行运行算子。

一、创建spring boot工程

pom.xml文件如下:

<?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 https://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <groupId>com.xxx</groupId>
    <artifactId>flatdemo</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <name>flatdemo</name>
    <description>FlatMap Demo</description>
    <properties>
        <java.version>1.8</java.version>
        <maven.compiler.source>8</maven.compiler.source>
        <maven.compiler.target>8</maven.compiler.target>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    </properties>
<dependencies>
    <dependency>
        <groupId>com.alibaba</groupId>
        <artifactId>fastjson</artifactId>
        <version>1.2.83</version>
    </dependency>
    <!-- Flink dependencies -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-core</artifactId>
        <version>1.13.5</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-java</artifactId>
        <version>1.13.5</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-java_2.11</artifactId>
        <version>1.13.5</version>
    </dependency>
</dependencies>

    <build>
        <plugins>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-assembly-plugin</artifactId>
                <version>2.5.5</version>
                <configuration>
                    <descriptorRefs>
                        <descriptorRef>jar-with-dependencies</descriptorRef>
                    </descriptorRefs>
                </configuration>
                <executions>
                    <execution>
                        <id>make-assembly</id>
                        <phase>package</phase>
                        <goals>
                            <goal>single</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>

</project>

二、map算子

map算子的功能:将一个数据项,通过map中的函数映射变为一个新的元素。
创建java类:

package com.xxx.flatdemo;

import org.apache.flink.api.common.functions.RichMapFunction;

public class MyMap extends RichMapFunction<String, String> {

    @Override
    public String map(String value)  {
        // 字符串转换为大写
        return value.toUpperCase();
    }
}

三、flatMap算子

flatMap算子的功能:将一个数据项,通过map中的函数映射生成零个、一个或者多个元素。
创建java类:

package com.xxx.flatdemo;

import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.util.Collector;

public class MyFlatMap extends RichFlatMapFunction<String,String> {


    @Override
    public void flatMap(String value, Collector<String> out) throws Exception {

        if(value!=null && value.startsWith("{") && value.endsWith("}")){
            JSONObject obj = JSON.parseObject(value);
            // 获取json的name字段
            String name = obj.getString("name");
            //按空格把名字拆散为name1和name2两个字段,如果名字只有一个单词,则name2为空字符串。
            String[] szName = name.split(" ");
            String name1;
            String name2;
            if (szName.length > 1) {
                name1 = szName[0];
                name2 = szName[1];
            }
            else {
                name1 = szName[0];
                name2 = "";
            }
            int age = obj.getIntValue("age");
            String tel = obj.getString("tel");
            String address = obj.getString("address");
            //拼接字段,字段之间用\u0001分隔,存储为Hive格式数据
            String line = name1+"\u0001" + name2 + "\u0001" + age + "\u0001" + tel + "\u0001" + address;
            out.collect(line);
        }
    }

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
    }
}

四、filter算子

filter算子的功能:对每个数据流中的每个元素执行一个布尔函数,只保留返回值为 True 的元素。
创建java类:

package com.xxx.flatdemo;

import org.apache.flink.api.common.functions.RichFilterFunction;

import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;

public class MyFilter extends RichFilterFunction<String> {

    @Override
    public boolean filter(String value) throws Exception {
        if(value!=null && value.startsWith("{") && value.endsWith("}")){
            JSONObject obj = JSON.parseObject(value);
            int age = obj.getIntValue("age");
            //过滤年龄大于等于100的记录
            if(age >= 100){
                return true;
            }else{
                return false;
            }
        } else {
            return false;
        }
    }
}
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容