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

Flink v1.13中ProcessStream算子高繁忙时间问题排查与解决

Flink应用高负载问题定位与解决方案

经过一个月的排查,最终定位到核心问题:

  • IP2Location库性能瓶颈:使用IP2Location Java库查询BIN文件中的IP地址位置时,峰值时段会引发性能问题。通过在读取BIN文件前传入IP2Proxy.IOModes.IP2PROXY_MEMORY_MAPPED参数,可避免该问题。
  • 状态对象不符合POJO标准:部分状态对象未遵循POJO规范,导致高负载。

环境与症状

  • 使用Flink v1.13版本,集群包含4个TaskManager(每个16核),共3800个任务,应用默认并行度为28。
  • 应用中ProcessStream算子长期处于高繁忙状态(占比约80%-90%),重启应用后繁忙时间暂时下降,但运行5-10小时后会再次上升。
  • 通过Grafana监控可见ProcessStream的繁忙时间持续增加,使用的Prometheus查询语句:avg((avg_over_time(flink_taskmanager_job_task_busyTimeMsPerSecond[1m]))) by (task_name)。
  • 该任务无背压,已通过flink_taskmanager_job_task_backPressuredTimeMsPerSecond指标验证。

相关代码

private void processOne(DataStream<KafkaObject> kafkaLog) {
    kafkaLog
         .filter(new FilterRequest())
         .name(FilterRequest.class.getSimpleName())
         .map(new MapToUserIdAndTimeStampMs())
         .name(MapToUserIdAndTimeStampMs.class.getSimpleName())
         .keyBy(UserObject::getUserId) // returns of type int
         .process(new ProcessStream())
         .name(ProcessStream.class.getSimpleName())
         .addSink(...)
         
        ;
}

// ...
// ...

public class ProcessStream extends KeyedProcessFunction<Integer, UserObject, Output>
{
    private static final long STATE_TIMER = // 5 min in milliseconds;

    private static final int AVERAGE_REQUEST = 74;
    private static final int STANDARD_DEVIATION = 32;
    private static final int MINIMUM_REQUEST = 50;
    private static final int THRESHOLD = 70;


    private transient ValueState<Tuple2<Integer, Integer>> state;

    @Override
    public void open(Configuration parameters) throws Exception
    {
        ValueStateDescriptor<Tuple2<Integer, Integer>> stateDescriptor = new ValueStateDescriptor<Tuple2<Integer, Integer>>(
                ProcessStream.class.getSimpleName(),
                TypeInformation.of(new TypeHint<Tuple2<Integer, Integer>>() {}));

        state = getRuntimeContext().getState(stateDescriptor);
    }

    @Override
    public void processElement(UserObject value, KeyedProcessFunction<Integer, UserObject, Output>.Context ctx, Collector<Output> out) throws Exception
    {
        Tuple2<Integer, Integer> stateValue = state.value();

        if (Objects.isNull(stateValue)) {
            stateValue = Tuple2.of(1, 0);
            ctx.timerService().registerProcessingTimeTimer(value.getTimestampMs() + STATE_TIMER);
        }

        int totalRequest = stateValue.f0;
        int currentScore = stateValue.f1;

        if (totalRequest >= MINIMUM_REQUEST && currentScore >= THRESHOLD)
        {
            out.collect(convert_to_output);
            state.clear();
        }
        else
        {
            stateValue.f0 = totalRequest + 1;
            stateValue.f1 =  calculateNextScore(stateValue.f0);
            state.update(stateValue);
        }
    }

    private int calculateNextScore(int totalRequest)
    {
        return (totalRequest - AVERAGE_REQUEST ) / STANDARD_DEVIATION;
    }

    @Override
    public void onTimer(long timestamp, KeyedProcessFunction<Integer, UserObject, Output>.OnTimerContext ctx, Collector<Output> out) throws Exception
    {
        state.clear();
    }
}

内容的提问来源于stack exchange,提问作者monstereo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 14:25:17