You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 06:57:24