Flink1.13.6提交任务报shuffle descriptors非法状态异常求助
问题背景
- 集群运行组件版本:Apache Flink 1.13.6、Scala 2.11、Kafka 2.2.2
- 作业逻辑:从Kafka读取逗号分隔的电影评分数据,完成字段拆分、多维度统计后,将统计结果写入Redis
- 故障现象:提交作业后任务启动失败,抛出非法状态异常
完整异常栈
Caused by: java.util.concurrent.CompletionException: java.lang.IllegalStateException: Trying to work with offloaded serialized shuffle descriptors. at java.util.concurrent.CompletableFuture.encodeRelay(CompletableFuture.java:326) at java.util.concurrent.CompletableFuture.completeRelay(CompletableFuture.java:338) at java.util.concurrent.CompletableFuture.uniRelay(CompletableFuture.java:925) at java.util.concurrent.CompletableFuture$UniRelay.tryFire(CompletableFuture.java:913) at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:488) at java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1990) at org.apache.flink.runtime.rpc.akka.AkkaInvocationHandler.lambda$invokeRpc$0(AkkaInvocationHandler.java:234) at java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:774) at java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:750) at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:488) at java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1990) at org.apache.flink.runtime.concurrent.FutureUtils$1.onComplete(FutureUtils.java:1079) at akka.dispatch.OnComplete.internal(Future.scala:263) at akka.dispatch.OnComplete.internal(Future.scala:261) at akka.dispatch.japi$CallbackBridge.apply(Future.scala:191) at akka.dispatch.japi$CallbackBridge.apply(Future.scala:188) at scala.concurrent.impl.CallbackRunnable.run(Promise.scala:36) at org.apache.flink.runtime.concurrent.Executors$DirectExecutionContext.execute(Executors.java:73) at scala.concurrent.impl.CallbackRunnable.executeWithValue(Promise.scala:44) at scala.concurrent.impl.Promise$DefaultPromise.tryComplete(Promise.scala:252) at akka.pattern.PromiseActorRef.$bang(AskSupport.scala:572) at akka.remote.DefaultMessageDispatcher.dispatch(Endpoint.scala:101) at akka.remote.EndpointReader$$anonfun$receive$2.applyOrElse(Endpoint.scala:999) at akka.actor.Actor$class.aroundReceive(Actor.scala:517) at akka.remote.EndpointActor.aroundReceive(Endpoint.scala:458) ... 9 more Caused by: java.lang.IllegalStateException: Trying to work with offloaded serialized shuffle descriptors. at org.apache.flink.runtime.deployment.InputGateDeploymentDescriptor.getShuffleDescriptors(InputGateDeploymentDescriptor.java:150) at org.apache.flink.runtime.io.network.partition.consumer.SingleInputGateFactory.create(SingleInputGateFactory.java:125) at org.apache.flink.runtime.io.network.NettyShuffleEnvironment.createInputGates(NettyShuffleEnvironment.java:261) at org.apache.flink.runtime.taskmanager.Task.<init>(Task.java:420) at org.apache.flink.runtime.taskexecutor.TaskExecutor.submitTask(TaskExecutor.java:737) at sun.reflect.GeneratedMethodAccessor32.invoke(Unknown Source) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.lambda$handleRpcInvocation$1(AkkaRpcActor.java:316) at org.apache.flink.runtime.concurrent.akka.ClassLoadingUtils.runWithContextClassLoader(ClassLoadingUtils.java:83) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcInvocation(AkkaRpcActor.java:314) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:217) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:163) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:24) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:20) at scala.PartialFunction.applyOrElse(PartialFunction.scala:123) at scala.PartialFunction.applyOrElse$(PartialFunction.scala:122) at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:20) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:172) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:172) at akka.actor.Actor.aroundReceive(Actor.scala:537) at akka.actor.Actor.aroundReceive$(Actor.scala:535) at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:220) at akka.actor.ActorCell.receiveMessage(ActorCell.scala:580) at akka.actor.ActorCell.invoke(ActorCell.scala:548) at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:270) at akka.dispatch.Mailbox.run(Mailbox.scala:231) at akka.dispatch.Mailbox.exec(Mailbox.scala:243) at java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:289) at java.util.concurrent.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1056) at java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1692) at java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:175)
作业核心代码
object batchProcess { def main(args:Array[String]): Unit = { val host = "localhost" val port = 6379 val env = StreamExecutionEnvironment.getExecutionEnvironment // read from kafka val source = KafkaSource.builder[String].setBootstrapServers("localhost:9092") .setTopics("movie_rating_records").setGroupId("my-group").setStartingOffsets(OffsetsInitializer.earliest) .setValueOnlyDeserializer(new SimpleStringSchema()) .setBounded(OffsetsInitializer.latest).build() // val inputDataStream = env.readTextFile("a.txt") val inputDataStream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source") val dataStream = inputDataStream .map( data =>{ val arr = data.split(",") ( arr(0),arr(1).toInt,arr(2).toInt,arr(3).toFloat,arr(4).toLong) }) val (counterUserIdPos,counterUserIdNeg,counterMovieIdPos,counterMovieIdNeg,counterUserId2MovieId) = commonProcess(dataStream) counterUserIdPos.map(x =>{ val jedisIns = new Jedis(host,port,100000) jedisIns.set("batch2feature_userId_rating1_"+x._1.toString, x._2.toString) jedisIns.close() }) env.execute("test") } }
排查方向与解决思路
优先排查shuffle描述符配置不一致问题
该异常是Flink 1.13版本典型的配置不匹配问题:Flink从1.13开始支持将大体积的shuffle部署描述符从RPC消息卸载到BlobServer存储,避免RPC消息体积过大触发帧大小超限,如果JobManager和TaskManager加载的该配置不一致,就会出现JM按卸载逻辑序列化描述符、TM按非卸载逻辑读取直接抛错。
- 检查集群所有节点的
flink-conf.yaml配置是否完全同步,确认jobmanager.deployment.shuffle-descriptors.offload.enabled配置在所有JM、TM节点取值一致 - 快速验证可直接在所有节点的配置文件中显式添加配置
jobmanager.deployment.shuffle-descriptors.offload.enabled: false,同步配置后重启集群重新提交作业 - 检查提交作业时是否通过
-D参数单独给JM设置了该配置,导致TM未加载到相同参数值
排查依赖冲突问题
如果配置确认一致仍报错,需要检查作业打包的依赖是否存在版本冲突:
- 打包作业时,所有Flink内核相关依赖(
flink-runtime、flink-core、flink-streaming-java等)、Scala标准库依赖必须设置为provided范围,不要打入作业fatjar,避免TM加载作业内自带的、和集群版本不匹配的类,导致序列化/反序列化逻辑不兼容 - 检查集群
lib目录下是否存在混入的其他版本Flink依赖包,导致类加载时加载了错误版本的类
修复作业代码中存在的其他问题
即使解决启动报错,当前代码也存在会导致运行失败的问题,需要同步修复:
- Kafka、Redis连接地址硬编码为
localhost,作业提交到集群运行时会连接所在TaskManager节点的本地端口,无法访问实际部署的服务,需要替换为集群内可路由的实际服务地址 - Redis写入逻辑在map算子中每条数据创建一次Jedis连接,会频繁创建销毁TCP连接,性能极差还会打满Redis连接数,需要改用
RichMapFunction在算子open方法中初始化连接、close方法中释放连接,或直接使用官方RedisSink连接器实现写入 - 示例数据首字段为
1542295208rating,拆分后转Int/Float/Long类型时会抛出数字格式异常,需要提前做字段格式校验或修正数据拆分逻辑
内容的提问来源于stack exchange,提问作者Sam Chiu
相关产品推荐
相关产品推荐

