可以使用 Flink 的 filter 算子过滤掉非 JSON 格式的字符串,然后使用 Flink 的 map 算子提取 ts 字段作为时间戳,并使用 Flink 的 assignTimestampsAndWatermarks 算子生成水位线。

下面是一个示例代码:

DataStream<String> logs = ...; // 日志数据流

DataStream<String> filteredLogs = logs.filter(new FilterFunction<String>() {
    @Override
    public boolean filter(String value) throws Exception {
        try {
            // 尝试解析为 JSON 格式
            JsonNode jsonNode = new ObjectMapper().readTree(value);
            return true;
        } catch (Exception e) {
            // 解析失败,说明是非 JSON 格式的字符串,过滤掉
            return false;
        }
    }
});

DataStream<Tuple2<Long, String>> timestampedLogs = filteredLogs.map(new MapFunction<String, Tuple2<Long, String>>() {
    @Override
    public Tuple2<Long, String> map(String value) throws Exception {
        JsonNode jsonNode = new ObjectMapper().readTree(value);
        long ts = jsonNode.get('ts').asLong();
        return Tuple2.of(ts, value);
    }
}).assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor<Tuple2<Long, String>>(Time.seconds(5)) {
    @Override
    public long extractTimestamp(Tuple2<Long, String> element) {
        return element.f0;
    }
});

在上面的代码中,首先使用 filter 算子过滤掉非 JSON 格式的字符串,然后使用 map 算子提取 ts 字段作为时间戳,并生成 Tuple2<Long, String> 类型的数据流。最后使用 assignTimestampsAndWatermarks 算子生成水位线,水位线的最大乱序程度为 5 秒。

Flink 实战:过滤日志脏数据并提取时间戳生成水位线

原文地址: https://www.cveoy.top/t/topic/f1o8 著作权归作者所有。请勿转载和采集!

免费AI点我,无需注册和登录