如何使用Flink按时间过滤提取指定时间窗口的时序数据
Flink实现指定时间点时序数据筛选方案
实现思路
- 原始数据存在重复记录,先按
id+事件时间做去重,避免重复结果 - 采用数据自带的
time字段作为事件时间,配置合理Watermark容忍少量乱序数据 - 按id分组后通过键控状态存储各时间点的数值,注册目标时间点的事件时间定时器,触发计算时筛选出1小时前、45分钟前、30分钟前、15分钟前、当前时间共5个时间点的记录,按指定格式输出
- 样例中参考触发时间为
2022-06-28 16:00:00,如果是长期运行的定时计算场景,只需修改定时器注册逻辑为每15分钟触发一次即可。
Java实现代码片段
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner; import org.apache.flink.api.common.eventtime.WatermarkStrategy; 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.api.java.tuple.Tuple3; import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.KeyedProcessFunction; import org.apache.flink.util.Collector; import java.time.Duration; import java.time.LocalDateTime; import java.time.ZoneOffset; import java.time.format.DateTimeFormatter; import java.util.Arrays; import java.util.List; public class TimeSeriesFilter { // 时间格式定义 private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); // 触发计算的参考时间 private static final LocalDateTime TARGET_TRIGGER_TIME = LocalDateTime.parse("2022-06-28 16:00:00", FORMATTER); // 需要提取的时间偏移量(单位:分钟):1小时前、45分钟前、30分钟前、15分钟前、当前 private static final List<Integer> TIME_OFFSETS = Arrays.asList(60, 45, 30, 15, 0); public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // 替换为实际数据源,此处为模拟输入 DataStream<Tuple3<String, String, String>> sourceStream = env.fromElements( Tuple3.of("a1", "2022-06-28 15:00:00", "0.23"), Tuple3.of("a1", "2022-06-28 15:15:00", "0.89"), Tuple3.of("a1", "2022-06-28 15:30:00", "0.12"), Tuple3.of("a1", "2022-06-28 15:45:00", "0.45"), Tuple3.of("a1", "2022-06-28 16:00:00", "0.11"), Tuple3.of("b1", "2022-06-28 15:00:00", "0.23"), Tuple3.of("b1", "2022-06-28 15:15:00", "0.89"), Tuple3.of("b1", "2022-06-28 15:30:00", "0.34"), Tuple3.of("b1", "2022-06-28 15:45:00", "0.56"), Tuple3.of("b1", "2022-06-28 16:00:00", "0.11"), Tuple3.of("c1", "2022-06-28 15:00:00", "0.23"), Tuple3.of("c1", "2022-06-28 15:15:00", "0.89"), Tuple3.of("c1", "2022-06-28 15:30:00", "0.78"), Tuple3.of("c1", "2022-06-28 15:45:00", "0.90"), Tuple3.of("c1", "2022-06-28 16:00:00", "0.11"), // 模拟重复数据 Tuple3.of("a1", "2022-06-28 16:00:00", "0.11") ); // 分配事件时间戳与Watermark,容忍10秒乱序 DataStream<Tuple3<String, String, String>> timedStream = sourceStream.assignTimestampsAndWatermarks( WatermarkStrategy.<Tuple3<String, String, String>>forBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((SerializableTimestampAssigner<Tuple3<String, String, String>>) (element, recordTimestamp) -> { LocalDateTime time = LocalDateTime.parse(element.f1, FORMATTER); return time.toInstant(ZoneOffset.ofHours(8)).toEpochMilli(); }) ); // 按id分组,处理去重、状态存储、定时输出 DataStream<String> resultStream = timedStream.keyBy(data -> data.f0) .process(new KeyedProcessFunction<String, Tuple3<String, String, String>, String>() { // 存储每个时间点对应的值,key为时间戳(毫秒) private transient MapState<Long, String> timeValueState; // 去重标记状态 private transient ValueState<Boolean> existFlag; @Override public void open(Configuration parameters) { timeValueState = getRuntimeContext().getMapState(new MapStateDescriptor<>( "time-value", Long.class, String.class )); existFlag = getRuntimeContext().getState(new ValueStateDescriptor<>( "exist-flag", Boolean.class )); } @Override public void processElement(Tuple3<String, String, String> value, Context ctx, Collector<String> out) throws Exception { long eventTs = LocalDateTime.parse(value.f1, FORMATTER).toInstant(ZoneOffset.ofHours(8)).toEpochMilli(); // 重复数据直接过滤 if (Boolean.TRUE.equals(existFlag.value())) { return; } existFlag.update(true); timeValueState.put(eventTs, value.f2); // 注册目标触发时间的定时器,仅注册一次 long triggerTs = TARGET_TRIGGER_TIME.toInstant(ZoneOffset.ofHours(8)).toEpochMilli(); ctx.timerService().registerEventTimeTimer(triggerTs); } @Override public void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) throws Exception { String currentId = ctx.getCurrentKey(); long triggerTs = TARGET_TRIGGER_TIME.toInstant(ZoneOffset.ofHours(8)).toEpochMilli(); // 遍历需要输出的时间偏移量,按格式输出 for (Integer offsetMin : TIME_OFFSETS) { long targetTs = triggerTs - offsetMin * 60 * 1000L; LocalDateTime targetTime = LocalDateTime.ofInstant(java.time.Instant.ofEpochMilli(targetTs), ZoneOffset.ofHours(8)); // 样例输出value固定为1,可根据实际业务替换为timeValueState.get(targetTs)取真实值 out.collect(String.format("(\"%s\", \"%s\", \"1\")", currentId, targetTime.format(FORMATTER))); } // 清空状态 timeValueState.clear(); existFlag.clear(); } }); resultStream.print(); env.execute("time-series-filter-job"); } }
关键逻辑说明
- 去重实现:通过
ValueState标记id+事件时间维度的数据是否已经处理过,重复流入的数据直接丢弃,不会进入状态存储 - Watermark配置:针对15分钟固定粒度的时序数据,设置10秒乱序容忍即可覆盖绝大多数数据延迟场景,不会导致输出延迟过高
- 状态设计:使用
MapState存储单id下所有时间点的数值,避免多值状态的序列化开销,触发计算后直接清空状态,不会造成状态持续膨胀 - 扩展适配:如果需要作业长期运行,只需修改定时器注册逻辑,每次
onTimer触发计算后,注册下一个15分钟整时间点的定时器即可,其余逻辑无需调整;如果某时间点数据缺失,可直接在输出逻辑中补默认值。
内容的提问来源于stack exchange,提问作者Criwran
相关产品推荐
相关产品推荐

