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

Flink如何统计全局DOWN状态设备数及各区域DOWN状态设备数量

代码笔误修正

先修正你现有代码中的几个笔误,避免运行报错:

  1. genericStream.keyBy(GenericEvent::getDevice) 应为 GenericEvent::getDeviceName,匹配POJO定义的字段名
  2. 过滤逻辑中 deviceState.getState() 应为 deviceState.getStatus(),匹配DeviceState的字段名
  3. SingleOutputStreamOperator<DeviceState> downNodes = nodesStatuses.filter(nodeDownFilter) 中的nodesStatuses应为之前定义的deviceStateStream

前置依赖修改

要实现按区域统计,需要把原始事件的区域属性传递到设备状态流中,首先补充DeviceState的字段定义:

class DeviceState{
    String deviceName, status, region;
    Long timestamp;
    // 补充对应字段的getter、setter、全参构造方法
}

同时修改EventEvaluator的逻辑,生成DeviceState对象时,把原始GenericEvent的region字段赋值到DeviceState中。


需求1:统计全局状态为DOWN的设备总数量

默认按去重逻辑实现(同一设备多次上报DOWN状态不重复计数),如果需要统计DOWN事件的总触发次数可直接用sum聚合:

import org.apache.flink.api.common.state.MapState;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.api.common.state.ValueState;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
import org.apache.flink.util.Collector;
import org.apache.flink.api.java.tuple.Tuple2;

// 全局DOWN设备数统计
SingleOutputStreamOperator<Tuple2<String, Integer>> globalDownCount = downNodes
    // 所有数据用固定key分到同一个分区做全局聚合
    .keyBy(state -> "GLOBAL")
    .process(new KeyedProcessFunction<String, DeviceState, Tuple2<String, Integer>>() {
        // 存储已统计的DOWN设备ID,避免重复计数
        private MapState<String, Boolean> downDeviceSet;
        // 存储当前全局DOWN设备总数
        private ValueState<Integer> totalCount;

        @Override
        public void open(Configuration parameters) throws Exception {
            downDeviceSet = getRuntimeContext().getMapState(
                new MapStateDescriptor<>("down_device_set", String.class, Boolean.class)
            );
            totalCount = getRuntimeContext().getState(
                new ValueStateDescriptor<>("global_down_count", Integer.class, 0)
            );
        }

        @Override
        public void processElement(DeviceState state, Context ctx, Collector<Tuple2<String, Integer>> out) throws Exception {
            String deviceName = state.getDeviceName();
            // 仅未统计过的设备才计入总数
            if (!downDeviceSet.contains(deviceName)) {
                downDeviceSet.put(deviceName, true);
                int newCount = totalCount.value() + 1;
                totalCount.update(newCount);
                // 输出最新的全局总数
                out.collect(Tuple2.of("global_down_total", newCount));
            }
        }
    });

// 可根据需要输出到控制台、数据库或其他存储
globalDownCount.print("全局DOWN设备数");

需求2:按区域统计各区域内DOWN状态的设备数量

逻辑和全局统计基本一致,仅keyBy的字段改为region即可:

// 按区域统计DOWN设备数
SingleOutputStreamOperator<Tuple2<String, Integer>> regionDownCount = downNodes
    .keyBy(DeviceState::getRegion)
    .process(new KeyedProcessFunction<String, DeviceState, Tuple2<String, Integer>>() {
        private MapState<String, Boolean> regionDownDeviceSet;
        private ValueState<Integer> regionCount;

        @Override
        public void open(Configuration parameters) throws Exception {
            regionDownDeviceSet = getRuntimeContext().getMapState(
                new MapStateDescriptor<>("region_down_device_set", String.class, Boolean.class)
            );
            regionCount = getRuntimeContext().getState(
                new ValueStateDescriptor<>("region_down_count", Integer.class, 0)
            );
        }

        @Override
        public void processElement(DeviceState state, Context ctx, Collector<Tuple2<String, Integer>> out) throws Exception {
            String deviceName = state.getDeviceName();
            if (!regionDownDeviceSet.contains(deviceName)) {
                regionDownDeviceSet.put(deviceName, true);
                int newCount = regionCount.value() + 1;
                regionCount.update(newCount);
                // 输出格式:(区域名, 该区域DOWN设备数)
                out.collect(Tuple2.of(ctx.getCurrentKey(), newCount));
            }
        }
    });

regionDownCount.print("区域DOWN设备数");

补充说明

  • 如果需要支持设备从DOWN转为UP时自动扣减计数,可以把包含UP/DOWN全量状态的deviceStateStream接入处理逻辑,判断status为UP时从MapState中删除对应设备,同时扣减计数
  • 如果需要统计固定窗口内的指标,在keyBy之后添加对应的窗口分配器即可,比如统计每小时的DOWN设备数可以加window(TumblingProcessingTimeWindows.of(Time.hours(1)))
  • 如果设备量级极大,全局统计用固定key存在性能瓶颈,可以改为两阶段聚合:第一阶段先按设备哈希分片做局部去重统计,第二阶段再汇总全局总数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 01:48:04