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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 12:48:18