首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >无法使用Flink CLI将流部署到Apache Flink的HA集群

无法使用Flink CLI将流部署到Apache Flink的HA集群
EN

Stack Overflow用户
提问于 2016-04-14 14:10:55
回答 1查看 1.4K关注 0票数 1

我可以将流部署到Apache的独立安装(使用一个JobManager和多个TaskManagers),没有问题:

代码语言:javascript
复制
bin/flink run -m example-app-1.stag.local:6123 -d -p 4 my-flow-fat-jar.jar <flow parameters>

但是,当我运行相同的命令并部署到独立的HA集群时,这个命令会引发错误:

代码语言:javascript
复制
------------------------------------------------------------
 The program finished with the following exception:

org.apache.flink.client.program.ProgramInvocationException: The program execution failed: JobManager did not respond within 60000 milliseconds
    at org.apache.flink.client.program.Client.runDetached(Client.java:406)
    at org.apache.flink.client.program.Client.runDetached(Client.java:366)
    at org.apache.flink.client.program.DetachedEnvironment.finalizeExecute(DetachedEnvironment.java:75)
    at org.apache.flink.client.program.Client.runDetached(Client.java:278)
    at org.apache.flink.client.CliFrontend.executeProgramDetached(CliFrontend.java:844)
    at org.apache.flink.client.CliFrontend.run(CliFrontend.java:330)
    at org.apache.flink.client.CliFrontend.parseParameters(CliFrontend.java:1189)
    at org.apache.flink.client.CliFrontend.main(CliFrontend.java:1239)
Caused by: org.apache.flink.runtime.client.JobTimeoutException: JobManager did not respond within 60000 milliseconds
    at org.apache.flink.runtime.client.JobClient.submitJobDetached(JobClient.java:221)
    at org.apache.flink.client.program.Client.runDetached(Client.java:403)
    ... 7 more
Caused by: java.util.concurrent.TimeoutException: Futures timed out after [60000 milliseconds]
    at scala.concurrent.impl.Promise$DefaultPromise.ready(Promise.scala:219)
    at scala.concurrent.impl.Promise$DefaultPromise.result(Promise.scala:223)
    at scala.concurrent.Await$$anonfun$result$1.apply(package.scala:190)
    at scala.concurrent.BlockContext$DefaultBlockContext$.blockOn(BlockContext.scala:53)
    at scala.concurrent.Await$.result(package.scala:190)
    at scala.concurrent.Await.result(package.scala)
    at org.apache.flink.runtime.client.JobClient.submitJobDetached(JobClient.java:218)
    ... 8 more

活动职务管理器将下列错误写入日志:

代码语言:javascript
复制
2016-04-14 13:54:44,160 WARN  akka.remote.ReliableDeliverySupervisor                        - Association with remote system [akka.tcp://flink@127.0.0.1:62784] has failed, address is now gated for [5000] ms. Reason is: [Disassociated].
2016-04-14 13:54:46,299 WARN  org.apache.flink.runtime.jobmanager.JobManager                - Discard message LeaderSessionMessage(null,TriggerSavepoint(5de582462f334caee4733c60c6d69fd7)) because the expected leader session ID Some(72630119-fd0a-40e7-8372-45c93781e99f) did not equal the received leader session ID None.

所以,我不明白是什么导致了这样的错误?

如有需要,请告知我。

附注:

从Flink仪表板部署对独立的HA集群很好。当我仅通过Flink CLI进行部署时,就会出现这样的问题。

更新

我清除了动物园管理员,清除了Flink在磁盘上使用的目录,并重新部署了Flink独立的HA集群。然后,我尝试使用bin/flink run命令运行flow。如您所见,flink--jobmanager-0-example-app-1.stag.local.log).只写了一行关于问题的JobManager

所有JobManagers和TaskManagers都使用相同的flink-conf.yaml

代码语言:javascript
复制
jobmanager.heap.mb: 1024
jobmanager.web.port: 8081

taskmanager.data.port: 6121
taskmanager.heap.mb: 2048
taskmanager.numberOfTaskSlots: 4
taskmanager.memory.preallocate: false
taskmanager.tmp.dirs: /flink/data/task_manager

blob.server.port: 6130
blob.storage.directory: /flink/data/blob_storage

parallelism.default: 4

state.backend: filesystem
state.backend.fs.checkpointdir: s3a://example-flink/checkpoints

restart-strategy: none
restart-strategy.fixed-delay.attempts: 2
restart-strategy.fixed-delay.delay: 60s

recovery.mode: zookeeper
recovery.zookeeper.quorum: zookeeper-1.stag.local:2181,zookeeper-2.stag.local:2181,zookeeper-3.stag.local:2181
recovery.zookeeper.path.root: /example/flink
recovery.zookeeper.storageDir: s3a://example-flink/recovery
recovery.jobmanager.port: 6123

fs.hdfs.hadoopconf: /flink/conf

因此,似乎独立的HA集群配置正确。

更新2

FYI:我想安装这里所描述的独立HA集群。不是纱线HA簇。

更新3

下面是由bin/flink CLI:flink-username-client-hostname.local.log创建的日志。

EN

回答 1

Stack Overflow用户

回答已采纳

发布于 2016-04-14 14:21:02

在HA模式下启动Flink集群时,JobManager地址及其领导id被写入指定的ZooKeeper集群。为了与JobManager通信,你不仅要知道地址,还要知道它的领导地址。因此,您必须在由CLI读取的“flink- CLI. you”中指定以下参数。

代码语言:javascript
复制
recovery.mode: zookeeper
recovery.zookeeper.quorum: address of your cluster
recovery.zookeeper.path.root: ZK path you've started your cluster with

有了这些信息,客户端就知道在哪里可以找到ZooKeeper集群,在哪里可以找到JobManager地址及其领导id。

票数 2
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/36625742

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档