Apache Flink SourceFunction输出序列化值问题及非阈值数据展示
你的Flink滑动窗口输出与数据展示问题解决方案
针对你遇到的两个问题——控制台输出序列化对象格式、需要展示未超过阈值的vTemp数据,我给你整理了具体的解决方法:
1. 让POJO以可读格式输出
你现在看到的fkwtest.MachineData@577aa2db是因为MachineData类没有重写toString()方法,Java默认的toString方法只会输出类名加对象的哈希值。只需要给MachineData添加自定义的toString()方法,就能打印出对象的具体字段内容:
public class MachineData { // 你的原有字段 private int id; private String engineId; private Date date; private int vTemp; // 原有构造方法、getter/setter方法... // 重写toString方法 @Override public String toString() { return "MachineData{" + "id=" + id + ", engineId='" + engineId + '\'' + ", date=" + new SimpleDateFormat("yyyy-MM-dd").format(date) + ", vTemp=" + vTemp + '}'; } }
这样修改后,再打印MachineData对象时,就会显示每个字段的具体值了。
2. 展示未超过阈值的vTemp数据
你当前只处理了vTemp>100的告警数据,要同时展示正常数据(vTemp<=100),可以将原始数据流拆分为两个分支,分别处理告警和正常数据:
另外注意你之前的窗口代码缺少实际的处理逻辑(比如apply或process),所以这段窗口代码不会产生任何输出,我也帮你补充了窗口内的处理示例(比如统计窗口内的平均温度):
修改后的AppWindow代码:
public class AppWindow { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.registerType(MachineData.class); env.setParallelism(2); env.enableCheckpointing(1500, CheckpointingMode.EXACTLY_ONCE); env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime); DataStream<MachineData> stream = env.addSource(new DataCollector()); // 分支1:处理并打印告警数据(vTemp>100) stream.filter(reading -> reading.getvTemp() > 100) .map(reading -> "-- ALERT Reading Above Threshold!: " + reading) .print(); // 分支2:处理并打印正常数据(vTemp<=100) stream.filter(reading -> reading.getvTemp() <= 100) .map(reading -> "-- Normal Reading: " + reading) .print(); // 滑动窗口处理示例:按engineId分组,统计10秒滑动窗口内的平均温度 stream.keyBy(reading -> reading.getEngineId()) .window(SlidingTimeWindows.of(Time.seconds(10), Time.seconds(10))) .apply(new WindowFunction<MachineData, String, String, TimeWindow>() { @Override public void apply(String engineId, TimeWindow window, Iterable<MachineData> dataList, Collector<String> out) throws Exception { int totalTemp = 0; int count = 0; for (MachineData data : dataList) { totalTemp += data.getvTemp(); count++; } double avgTemp = count > 0 ? (double) totalTemp / count : 0; String windowInfo = String.format("Window [%s - %s]", new SimpleDateFormat("HH:mm:ss").format(window.getStart()), new SimpleDateFormat("HH:mm:ss").format(window.getEnd())); out.collect(String.format("Engine %s | %s | Avg Temp: %.1f", engineId, windowInfo, avgTemp)); } }) .print(); env.execute("Machine Temperature Monitoring"); } }
关于Kryo警告的说明
你控制台中出现的Kryo反射访问警告是Flink默认序列化机制导致的,这个警告不会影响程序功能。如果想消除它,可以升级到较新版本的Flink(已适配Java 9+的模块系统),或者在启动程序时添加JVM参数--illegal-access=warn调整警告级别。
内容的提问来源于stack exchange,提问作者Stimks
相关产品推荐
相关产品推荐

