前言
电力线上业务夜间割接,大数据平台实时计算Flink作业需要配合重启操作。批量重启遇到大量Flink作业启动失败的现象,关键日志如下:

文章发布于微信公众号:大数据从业者,其它均为转载,原创不易,欢迎您点赞关注推荐转发,谢谢!

问题分析
java.util.concurrent.CompletionException: org.apache.flink.util.FlinkRuntimeException:Could not retrieve JobResults of globally-terminated jobs from JobResultStoreat java.util.concurrent.CompletableFuture.encodeThrowable(CompletableFuture.java:273)at java.util.concurrent.CompletableFuture.completeThrowable(CompletableFuture.java:280)at java.util.concurrent.CompletableFuture$AsyncSupply.run(CompletableFuture.java:1606)at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)at java.lang.Thread.run(Thread.java:750)Caused by: org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.exc.MismatchedInputException:No content to map due to end-of-input at [Source: (org.apache.flink.runtime.fs.hdfs.HadoopDataInputStream); line: 1, column: 0]at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.exc.MismatchedInputException.from(MismatchedInputException.java:59)at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper._initForReading(ObjectMapper.java:4765)at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper._readMapAndClose(ObjectMapper.java:4667)at org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper.readValue(ObjectMapper.java:3666)at org.apache.flink.runtime.highavailability.FileSystemJobResultStore.getDirtyResultsInternal(FileSystemJobResultStore.java:208)at org.apache.flink.runtime.highavailability.AbstractThreadsafeJobResultStore.withReadLock(AbstractThreadsafeJobResultStore.java:118)at org.apache.flink.runtime.highavailability.AbstractThreadsafeJobResultStore.getDirtyResults(AbstractThreadsafeJobResultStore.java:100)at org.apache.flink.runtime.dispatcher.runner.SessionDispatcherLeaderProcess.getDirtyJobResults(SessionDispatcherLeaderProcess.java:194)at org.apache.flink.runtime.dispatcher.runner.AbstractDispatcherLeaderProcess.supplyUnsynchronizedIfRunning(AbstractDispatcherLeaderProcess.java:198)at org.apache.flink.runtime.dispatcher.runner.SessionDispatcherLeaderProcess.getDirtyJobResultsIfRunning(SessionDispatcherLeaderProcess.java:188)at java.util.concurrent.CompletableFuture$AsyncSupply.run(CompletableFuture.java:1604)
从上述日志来看异常发生于高可用Highavailability模块,而JobResultStore是Flink的一个组件,用于将全局终止(即已完成、已取消或失败)作业的结果持久化存储至文件系统,确保这些结果在作业完成后仍能保留,这些结果随后会被Flink用于判断作业是否应在开启高可用场景中进行恢复。对应源码如图所示:

通过getDirtyJobResultsIfRunning方法获取到Collection

查看该参数说明如下:

默认值为:
{high-availability.storageDir}/job-result-store/{high-availability.cluster-id}
而{high-availability.cluster-id}其实就是yarn applicationId(每个作业都唯一),那么新启动的作业的这个路径下应该是空目录,而不应该存在DIRTY.json。
随即查看启动失败的Flink作业参数job-result-store.storage-path配置值发现都是固定相同的路径:/flink/ha/job-result-store。
解决措施
将所有Flink作业job-result-store.storage-path不显示配置,即采用默认值。或者如果非要显示配置该参数,需要自行设置每个作业使用唯一独立的路径。
总结
至此,一个简单的配置引发的大规模Flink实时业务停服故障,排查处理完成!文章发布于微信公众号:大数据从业者,其它均为转载,原创不易,欢迎您点赞关注推荐转发,谢谢!





