IDEA配置flink开发环境及local集群代码测试
背景:最近公司需要引入flink相关框架做一些大数据报表分析的任务,之前没有实际接触过flink,所以需要学习一下。此外,防止看完就忘,也为了后续的回顾学习,因此在这里做一个整理,也希望帮助到有需要的朋友。环境准备:我这里是在自己的笔记本上搭建的环境VMware 安装centos7虚拟机 并配置好网络等win10安装idea 并配置maven(要求3.0以上,我用的3.6.2)...
背景:
最近公司需要引入flink相关框架做一些大数据报表分析的任务,之前没有实际接触过flink,所以需要学习一下。此外,防止看完就忘,也为了后续的回顾学习,因此在这里做一个整理,也希望帮助到有需要的朋友。
环境准备:
我这里是在自己的笔记本上搭建的环境
- VMware 安装centos7虚拟机 并配置好网络等
- win10安装idea 并配置maven(要求3.0以上,我用的3.6.2)
- flink-1.7.2-bin-hadoop27-scala_2.12.tgz
- jdk(要求1.8以上)
了解到有三种配置idea开发flink的方式
-
通过cmd运行mvn archetype:generate来生成flink模板项目,然后mvn clean package进行编译,之后导入idea即可开发
-
直接在idea中创建一个maven项目(可以是空的maven项目,也可以通过选择archetype创建模板项目,模板会将依赖生成好,不需要很大的修改),然后在pom中配置flink相关依赖即可开发
-
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项目:
更多推荐
所有评论(0)