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
相关产品推荐
相关产品推荐

