暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

参数配置错误:南方电力业务割接期间导致大量Flink作业重启失败问题排查

大数据从业者 2025-10-16
121

前言

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

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

问题分析

    java.util.concurrent.CompletionException: org.apache.flink.util.FlinkRuntimeException: 
      Could not retrieve JobResults of globally-terminated jobs from JobResultStore 
      at 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,如果不为空,则通过recoverJobsIfRunning方法的recoverJobs方法尝试恢复Job,以达到高可用的效果。看着流程也没啥问题,而异常是获取JobResult时候,提示解析JSON异常。通过查看HDFS确实存在DIRTY.json的空文件。而这个空文件是由配置参数job-result-store.storage-path决定的,对应源码如图所示: 

    查看该参数说明如下:

    默认值为:

      {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实时业务停服故障,排查处理完成!文章发布于微信公众号:大数据从业者,其它均为转载,原创不易,欢迎您点赞关注推荐转发,谢谢!

      文章转载自大数据从业者,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

      评论