Skip to main content

安装

安装包下载

链接:https://pan.baidu.com/s/1nBrOPbvxn-tVq1DBwEkkaw 
提取码:7b1i
--来自百度网盘超级会员V5的分享

解压

tar -zxvf flink-1.16.0-bin-scala_2.12.tgz -C ../module/

提交作用的不同方式

Standalone运行模式

修改配置文件

vim flink-conf.yaml

JobManager节点地址(公共部分)

jobmanager.rpc.address: 10.240.8.68
jobmanager.bind-host: 0.0.0.0
rest.address: 10.240.8.68
rest.bind-address: 0.0.0.0
taskmanager.numberOfTaskSlots: 3

修改workers

vi workers
10.240.8.68
10.240.8.229
10.240.8.185

修改masters

vi masters
10.240.8.68:8081

同步安装包多所有的集群

./xsync.sh /home/bigdata/module/flink-1.16.0

TaskManager节点地址.需要配置为当前机器名(不同的work节点配置的host不一样就行了)

taskmanager.bind-host: 0.0.0.0
taskmanager.host: 10.240.8.229

启动集群

./start-cluster.sh

启动效果

master  StandaloneSessionClusterEntrypoint
work TaskManagerRunner

Yarn运行模式

配置环境变量

那台机器要启动flink就配置那台就行了,不用所有的hadoop集群都要配置这个

sudo vim /etc/profile.d/my_env.sh

export HADOOP_HOME=/home/bigdata/module/hadoop-3.2.3
export PATH=$PATH:$HADOOP_HOME/bin
export PATH=$PATH:$HADOOP_HOME/sbin
export HADOOP_CONF_DIR=${HADOOP_HOME}/etc/hadoop
export HADOOP_CLASSPATH=`hadoop classpath`


source /etc/profile.d/my_env.sh

YARN Session模式

执行脚本命令向YARN集群申请资源,开启一个YARN会话,启动Flink集群。

bin/yarn-session.sh -nm test
  • -d:分离模式,如果你不想让Flink YARN客户端一直前台运行,可以使用这个参数,即使关掉当前对话窗口,YARN session也可以后台运行。
  • -jm(--jobManagerMemory):配置JobManager所需内存,默认单位MB。
  • -nm(--name):配置在YARN UI界面上显示的任务名。
  • -qu(--queue):指定YARN队列名。
  • -tm(--taskManager):配置每个TaskManager所使用内存。

YARN的会话模式不会把集群资源固定,同样是动态分配的。

向集群提交任务

bin/flink run -c com.xx.xx xxx.jar

比如官网测试程序

./bin/flink run ./examples/streaming/TopSpeedWindowing.jar

YARN Per-job模式

在YARN环境中,由于有了外部平台做资源调度,所以我们也可以直接向YARN提交一个单独的作业,从而启动一个Flink集群。

./bin/flink run -t yarn-per-job --detached xx.jar

比如官网测试用例

./bin/flink run -t yarn-per-job --detached ./examples/streaming/TopSpeedWindowing.jar
yarn application -kill application_1698659960145_0005

yarn logs -applicationId application_1698659960145_0005

注意:如果启动过程中报如下异常。

Exception in thread “Thread-5” java.lang.IllegalStateException: Trying to access closed classloader. Please check if you store classloaders directly or indirectly in static fields. If the stacktrace suggests that the leak occurs in a third party library and cannot be fixed immediately, you can disable this check with the configuration ‘classloader.check-leaked-classloader’.
at org.apache.flink.runtime.execution.librarycache.FlinkUserCodeClassLoaders

解决方法:

vim flink-conf.yaml

classloader.check-leaked-classloader: false

YARN Application模式

可以通过yarn.provided.lib.dirs配置选项指定位置,将flink的依赖上传到远程。

  • 上传flink的lib和plugins到HDFS上。
hadoop fs -mkdir /flink-dist
hadoop fs -put /home/bigdata/module/flink-1.16.0/lib/ /flink-dist
hadoop fs -put /home/bigdata/module/flink-1.16.0/plugins/ /flink-dist
  • 上传自己的jar包到HDFS。
hadoop fs -mkdir /flink-jars
hadoop fs -put ./examples/streaming/TopSpeedWindowing.jar /flink-jars
  • 提交作业
bin/flink run-application -t yarn-application   -Dyarn.provided.lib.dirs="hdfs://bigdatacluster/flink-dist" -c com.xxx  hdfs://bigdatacluster/flink-jars/xx.jar

比如

bin/flink run-application -t yarn-application   -Dyarn.provided.lib.dirs="hdfs://bigdatacluster/flink-dist" hdfs://bigdatacluster/flink-jars/TopSpeedWindowing.jar

这种方式下,flink本身的依赖和用户jar可以预先上传到HDFS,而不需要单独发送到集群,这就使得作业提交更加轻量了。

配置历史服务器

运行 Flink job 的集群一旦停止,只能去 yarn 或本地磁盘上查看日志,不再可以查看作业挂掉之前的运行的 Web UI,很难清楚知道作业在挂的那一刻到底发生了什么。

创建对应的目录

hadoop fs -mkdir -p /logs/flink-job

修改配置文件

vi /home/bigdata/module/flink-1.16.0/conf/flink-conf.yaml
jobmanager.archive.fs.dir: hdfs://bigdatacluster/logs/flink-job
historyserver.web.address: 0.0.0.0
historyserver.web.port: 8082
historyserver.archive.fs.dir: hdfs://bigdatacluster/logs/flink-job
historyserver.archive.fs.refresh-interval: 5000

启停历史服务器

bin/historyserver.sh start

提交一个作业

bin/flink run-application -t yarn-application   -Dyarn.provided.lib.dirs="hdfs://bigdatacluster/flink-dist" hdfs://bigdatacluster/flink-jars/TopSpeedWindowing.jar

访问

http://master1:8082

停止历史服务器

bin/historyserver.sh stop