背景:

最近公司需要引入flink相关框架做一些大数据报表分析的任务,之前没有实际接触过flink,所以需要学习一下。此外,防止看完就忘,也为了后续的回顾学习,因此在这里做一个整理,也希望帮助到有需要的朋友。

环境准备:

我这里是在自己的笔记本上搭建的环境

  1. VMware 安装centos7虚拟机 并配置好网络等
  2. win10安装idea 并配置maven(要求3.0以上,我用的3.6.2)
  3. flink-1.7.2-bin-hadoop27-scala_2.12.tgz
  4. jdk(要求1.8以上)

了解到有三种配置idea开发flink的方式

  1. 通过cmd运行mvn archetype:generate来生成flink模板项目,然后mvn clean package进行编译,之后导入idea即可开发

  2. 直接在idea中创建一个maven项目(可以是空的maven项目,也可以通过选择archetype创建模板项目,模板会将依赖生成好,不需要很大的修改),然后在pom中配置flink相关依赖即可开发

  3. curl url...获取模板项目,导入idea即可开发  

 我这里只试了1和2两种方式,都可以。这里以第2种空模板的方式为例,最后面附带一个通过arachetype创建的截图。

1.首先创建一个空的maven项目

2.然后修改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 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>asn</groupId>
    <artifactId>flinkLearn</artifactId>
    <version>1.0-SNAPSHOT</version>
    <dependencies>
        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-java -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-java</artifactId>
            <version>1.7.2</version>
            <scope>provided</scope>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-streaming-java -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-java_2.12</artifactId>
            <version>1.7.2</version>
            <scope>provided</scope>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-scala -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-scala_2.12</artifactId>
            <version>1.7.2</version>
            <scope>provided</scope>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-streaming-scala -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-scala_2.12</artifactId>
            <version>1.7.2</version>
            <scope>provided</scope>

        </dependency>
    </dependencies>
</project>

<scope>provided</scope>可以保证在打包的时候不会把依赖的jar包也打进去,避免跟flink集群中的包冲突。但是在本地run的时候可能需要注释掉。也就是说只有打包的时候才使用。

3.等带maven将相关依赖下载下来(这个时间有点长,半小时多)

4.依赖下载完之后就可以开发flink程序了,下面是两个wordcount程序BatchWordCountJava和SocketWindowWordCountJava

 

 

 不知道为什么,我在虚拟机上下载的netcat不能有效监听端口(nc -l 9000),Windows本地的flink程序会报连接超时。但是下载netcat到Windows上,使用Windows的cmd监听是可以的,具体可以参考这里。这里先测试批量计算。下面是完整代码

package wordCount;

import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.operators.AggregateOperator;
import org.apache.flink.api.java.operators.DataSource;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.util.Collector;

public class BatchWordCountJava {
    public static void main(String[] args) throws Exception {
        String inputPath = "/opt/testBatch.txt";
        String outPath = "/opt/output";
        //获取运行环境
        ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
        //获取数据源
        DataSource<String> text = env.readTextFile(inputPath);

        AggregateOperator<Tuple2<String, Integer>> sum = text.flatMap(new Tokenizer()).groupBy(0).sum(1);

        sum.writeAsCsv(outPath,"\n"," ").setParallelism(1);
        env.execute();
    }

    public static class Tokenizer implements FlatMapFunction<String, Tuple2<String,Integer>>{

        @Override
        public void flatMap(String s, Collector<Tuple2<String, Integer>> collector) throws Exception {
            String[] split = s.toLowerCase().split("\\s+");
            for (String word:split){
                if (word.length()>0){
                    collector.collect(new Tuple2<>(word,1));
                }
            }
        }
    }
}

可以现在本地测试之后再打包到服务器上测试,只需要注意打包的时候修改输入文件路径和输出结果路径。(代码中的路径是我服务器上的路径,各位需要自己修改)

5.如果本地没问题,接下来需要打包,打包的话需要在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>asn</groupId>
    <artifactId>flinkLearn</artifactId>
    <version>1.0-SNAPSHOT</version>
    <dependencies>
        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-java -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-java</artifactId>
            <version>1.7.2</version>
            <scope>provided</scope><!-- 这个主要是保证打包的时候不会把额外的依赖也打包,避免跟集群中的包冲突 -->
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-streaming-java -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-java_2.12</artifactId>
            <version>1.7.2</version>
            <scope>provided</scope>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-scala -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-scala_2.12</artifactId>
            <version>1.7.2</version>
            <scope>provided</scope>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-streaming-scala -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-scala_2.12</artifactId>
            <version>1.7.2</version>
            <scope>provided</scope>
        </dependency>
    </dependencies>
    <build>
        <plugins>
            <!-- 编译插件 -->
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-compiler-plugin</artifactId>
                <version>3.6.0</version>
                <configuration>
                    <source>1.8</source>
                    <target>1.8</target>
                    <encoding>UTF-8</encoding>
                </configuration>
            </plugin>
            <!-- scala编译插件 -->
            <plugin>
                <groupId>net.alchim31.maven</groupId>
                <artifactId>scala-maven-plugin</artifactId>
                <version>3.1.6</version>
                <configuration>
                    <scalaCompatVersion>2.12</scalaCompatVersion>
                    <scalaVersion>2.12</scalaVersion>
                    <encoding>UTF-8</encoding>
                </configuration>
                <executions>
                    <execution>
                        <id>compile-scala</id>
                        <phase>compile</phase>
                        <goals>
                            <goal>add-source</goal>
                            <goal>compile</goal>
                        </goals>
                    </execution>
                    <execution>
                        <id>test-compile-scala</id>
                        <phase>test-compile</phase>
                        <goals>
                            <goal>add-source</goal>
                            <goal>testCompile</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>
            <!-- 打jar包插件(会包含所有依赖) -->
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-assembly-plugin</artifactId>
                <version>2.6</version>
                <configuration>
                    <descriptorRefs>
                        <descriptorRef>jar-with-dependencies</descriptorRef>
                    </descriptorRefs>
                    <archive>
                        <manifest>
                            <!-- 这里可以设置jar包的入口类(可选) 如果不设置,可以在命令行run jar包的时候通过-c指定入口类-->
                            <mainClass>wordCount.BatchWordCountJava</mainClass>
                        </manifest>
                    </archive>
                </configuration>
                <executions>
                    <execution>
                        <id>make-assembly</id>
                        <phase>package</phase>
                        <goals>
                            <goal>single</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>
</project>

从cmd中进入当前项目所在目录,执行mvn clean package进行打包,之后会在target文件夹中生产两个jar包(一个带依赖)

选择第二个带依赖的jar包上传到centos中,之后就可以启动flink运行这个程序了(保证存在对应的文件路径)

[root@flink1 flink-1.7.2]# bin/flink run ../flinkLearn-1.0-SNAPSHOT-jar-with-dependencies.jar

//显式指定入口类
//bin/flink run -c xxx ../flinkLearn-1.0-SNAPSHOT-jar-with-dependencies.jar

之后就能看到有output文件生成,且通过webui也可以看到任务执行情况。

(flink集群启动主要有三种模式,一种是直接start-cluster.sh启动本地模式,一种是standlone模式,还有一种就是yarn模式)

通过arachetype创建maven项目:

Logo

华为开发者空间,是为全球开发者打造的专属开发空间,汇聚了华为优质开发资源及工具,致力于让每一位开发者拥有一台云主机,基于华为根生态开发、创新。

更多推荐