Flink如何统计全局DOWN状态设备数及各区域DOWN状态设备数量
代码笔误修正
先修正你现有代码中的几个笔误,避免运行报错:
genericStream.keyBy(GenericEvent::getDevice)应为GenericEvent::getDeviceName,匹配POJO定义的字段名- 过滤逻辑中
deviceState.getState()应为deviceState.getStatus(),匹配DeviceState的字段名 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
相关产品推荐
相关产品推荐

