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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:48:11