Flink 1.17中RocksDB状态TTL触发时TaskManager失联故障排查
问题描述
我使用的是flink 1.17.0版本,有两个从Kafka消费的消息流A和B,数据量非常大。我需要将这两个流存储到rocksdb中进行关联处理,程序运行正常,但每次触发rocksdb的TTL清理时都会出现如下错误,即使我将TTL设置为2分钟(2分钟内的数据量约几百万条,单条数据大小为10kb)。我尝试过多种rocksdb配置组合,但问题仍未解决。
代码实现
public class AfterKeyByStateTTL extends KeyedCoProcessFunction<String, JSONObject, JSONObject, String> { private transient ValueState<JSONObject> newAState; private transient ListState<JSONObject> newBState; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); Map<String, String> globalMap = getRuntimeContext().getExecutionConfig().getGlobalJobParameters() .toMap(); boolean cleanupInRocksdbCompactFilter = Boolean.parseBoolean(globalMap.getOrDefault("app.join.state.ttl.cleanupInRocksdbCompactFilter", "false")); boolean cleanupFullSnapshot = Boolean.parseBoolean(globalMap.getOrDefault("app.join.state.ttl.cleanupFullSnapshot", "false")); boolean disableCleanupInBackground = Boolean.parseBoolean(globalMap.getOrDefault("app.join.state.ttl.disableCleanupInBackground", "false")); boolean neverReturnExpired = Boolean.parseBoolean(globalMap.getOrDefault("app.join.state.ttl.neverReturnExpired", "true")); ValueStateDescriptor<JSONObject> newADesc = new ValueStateDescriptor<>("newADesc", JSONObject.class); long ft = Long.parseLong(globalMap.getOrDefault("app.join.state.ttl.A-minutes", "1440")); Time fTime = Time.minutes(ft); long fNum = Long.parseLong(globalMap.getOrDefault("app.join.state.ttl.A-num", "1000")); StateTtlConfig.Builder fStateTtlConfig = StateTtlConfig.newBuilder(fTime) .updateTtlOnReadAndWrite() .cleanupIncrementally(10000, true); if (cleanupInRocksdbCompactFilter) fStateTtlConfig.cleanupInRocksdbCompactFilter(fNum); if (cleanupFullSnapshot) fStateTtlConfig.cleanupFullSnapshot(); if (disableCleanupInBackground) fStateTtlConfig.disableCleanupInBackground(); if (neverReturnExpired) fStateTtlConfig.neverReturnExpired(); newADesc.enableTimeToLive(fStateTtlConfig.build()); newAState = getRuntimeContext().getState(newADesc); ListStateDescriptor<JSONObject> newBDesc = new ListStateDescriptor<>("newBDesc", JSONObject.class); long pt = Long.parseLong(globalMap.getOrDefault("app.join.state.ttl.B-minutes", "1440")); Time pTime = Time.minutes(pt); long pNum = Long.parseLong(globalMap.getOrDefault("app.join.state.ttl.B-num", "1000")); StateTtlConfig.Builder pStateTtlConfig = StateTtlConfig.newBuilder(pTime) .updateTtlOnReadAndWrite() .cleanupIncrementally(10000, true); if (cleanupInRocksdbCompactFilter) pStateTtlConfig.cleanupInRocksdbCompactFilter(pNum); if (cleanupFullSnapshot) pStateTtlConfig.cleanupFullSnapshot(); if (disableCleanupInBackground) pStateTtlConfig.disableCleanupInBackground(); if (neverReturnExpired) pStateTtlConfig.neverReturnExpired(); newBDesc.enableTimeToLive(pStateTtlConfig.build()); newBState = getRuntimeContext().getListState(newBDesc); } @Override public void processElement1(JSONObject AValue, Context ctx, Collector<String> out) throws Exception { // joinKey is to get the associated key String AJoinKey = joinKey(AValue); Iterable<JSONObject> ite = newBState.get(); if (ite != null && ite.iterator().hasNext()) { for (JSONObject BValue : ite) { String BJoinKey = joinKey(BValue); if (AJoinKey.equals(BJoinKey)) { out.collect(Tools.getResult(AValue, BValue)); } } } newAState.update(AValue); } @Override public void processElement2(JSONObject BValue, Context ctx, Collector<String> out) throws Exception { JSONObject nv = newAState.value(); if (nv != null && joinKey(nv).equals(joinKey(BValue))) { out.collect(getResult(nv, BValue)); } else { newBState.add(BValue); } } @Override public void close() throws Exception { super.close(); }
错误日志
org.apache.flink.runtime.jobmaster.JobMasterException: TaskManager with id container_1747647448650_0425_01_000003(bd1:45454) is no longer reachable. at org.apache.flink.runtime.jobmaster.JobMaster$TaskManagerHeartbeatListener.notifyTargetUnreachable(JobMaster.java:1449) at org.apache.flink.runtime.heartbeat.DefaultHeartbeatMonitor.reportHeartbeatRpcFailure(DefaultHeartbeatMonitor.java:126) at org.apache.flink.runtime.heartbeat.HeartbeatManagerImpl.runIfHeartbeatMonitorExists(HeartbeatManagerImpl.java:275) at org.apache.flink.runtime.heartbeat.HeartbeatManagerImpl.reportHeartbeatTargetUnreachable(HeartbeatManagerImpl.java:267) at org.apache.flink.runtime.heartbeat.HeartbeatManagerImpl.handleHeartbeatRpcFailure(HeartbeatManagerImpl.java:262) at org.apache.flink.runtime.heartbeat.HeartbeatManagerImpl.lambda$handleHeartbeatRpc$0(HeartbeatManagerImpl.java:248) at java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:774) at java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:750) at java.util.concurrent.CompletableFuture$Completion.run(CompletableFuture.java:456) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.lambda$handleRunAsync$4(AkkaRpcActor.java:453) at org.apache.flink.runtime.concurrent.akka.ClassLoadingUtils.runWithContextClassLoader(ClassLoadingUtils.java:68) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRunAsync(AkkaRpcActor.java:453) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:218) at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:84) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:168) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:24) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:20) at scala.PartialFunction.applyOrElse(PartialFunction.scala:127) at scala.PartialFunction.applyOrElse$(PartialFunction.scala:126) at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:20) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:175) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176) 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:579) at akka.actor.ActorCell.invoke(ActorCell.scala:547) 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:1067) at java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1703) at java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:172) org.apache.flink.runtime.rpc.exceptions.RecipientUnreachableException: Could not send message [RemoteRpcInvocation(TaskExecutorGateway.submitTask(TaskDeploymentDescriptor, JobMasterId, Time))] from sender [Actor[akka://flink/temp/taskmanager_0$8qb]] to recipient [Actor[akka.tcp://flink@bd2:46648/user/rpc/taskmanager_0#379917460]], because the recipient is unreachable. This can either mean that the recipient has been terminated or that the remote RpcService is currently not reachable. at org.apache.flink.runtime.rpc.akka.DeadLettersActor.handleDeadLetter(DeadLettersActor.java:61) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:24) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:20) at scala.PartialFunction.applyOrElse(PartialFunction.scala:127) at scala.PartialFunction.applyOrElse$(PartialFunction.scala:126) at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:20) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:175) 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:579) at akka.actor.ActorCell.invoke(ActorCell.scala:547) 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:1067) at java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1703) at java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:172)
请问这是否是我的配置问题?
问题分析与解决建议
从错误日志和场景来看,TaskManager失联是TTL清理过程中资源过载导致进程崩溃或心跳超时,和你的配置直接相关,具体调整方向如下:
1. 降低增量清理的资源消耗
你设置的cleanupIncrementally(10000, true)会让业务处理线程每次处理记录时同步清理10000条过期状态,在高吞吐场景下会直接挤占业务处理的CPU资源,导致TaskManager无法及时响应JobMaster的心跳,被判定为失联。
调整建议:
- 把单次清理条数降到
1000以内,比如cleanupIncrementally(500, false); - 将第二个参数设为
false,让清理工作在后台线程执行,避免阻塞业务处理。
2. 关闭RocksDB压缩阶段的TTL清理
如果开启了cleanupInRocksdbCompactFilter,RocksDB会在压缩(Compaction)时清理过期状态。但你的状态量极大(2分钟几十GB),压缩本身就会占用大量CPU和IO资源,叠加TTL清理会直接压垮TaskManager。
调整建议:
- 禁用
cleanupInRocksdbCompactFilter,优先依赖增量清理和后台异步清理,Flink官方不建议在大状态场景下开启该选项。
3. 优化RocksDB内存配置
状态量巨大时,RocksDB的内存配置不足会导致频繁磁盘IO和写放大,进一步加剧资源消耗。
调整建议:
- 调大block cache大小,比如设置
state.backend.rocksdb.block.cache-size为TaskManager内存的30%-40%; - 增加write buffer大小,比如
state.backend.rocksdb.write-buffer-size设为64MB,减少压缩频率; - 检查TaskManager的JVM参数,确保堆内存和直接内存足够(RocksDB会使用直接内存存储缓存数据)。
4. 其他优化点
- 调整
neverReturnExpired参数:若业务允许,可设为false,减少状态访问时的TTL检查开销; - 优化状态结构:
newBState作为ListState,若单个key下数据过多,遍历和清理都会消耗大量资源,可考虑细化key的粒度,分散状态量; - 监控资源指标:开启Flink Metrics,观察TTL清理时的CPU、内存、IO使用率,确认资源瓶颈点。
内容的提问来源于stack exchange,提问作者datagic
相关产品推荐
相关产品推荐

