Flink 实战:过滤日志脏数据并提取时间戳生成水位线
可以使用 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 秒。
原文地址: https://www.cveoy.top/t/topic/f1o8 著作权归作者所有。请勿转载和采集!