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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 16:22:32