Skip to content
Charles Shao
Go back

Flink 深挖 · DataStream API:ProcessFunction、状态、Timer 与 Side Output

views

前面几篇里,无论是时间窗口还是端到端 Exactly-Once 的实时宽表,能用 Flink SQL 表达的我们都用 SQL——它声明式、上手快、优化器还帮你把状态和窗口管好。但真到广告实时链路的深水区,你会撞到一类 SQL 写不出来的需求:“某个广告主 15 分钟内没完成转化就告警”、“这条曝光迟到了 3 秒,别丢,单独送去补算”、“按用户维护一份随时能读写的实时画像”——这些要的是对状态、时间、分流的细粒度控制,声明式 SQL 给不了。

这时候就得下沉到 Flink 的底层 API——DataStream。这一篇承接前面的 SQL 与真实数据源,把 DataStream 最核心的一套武器讲透并跑通:KeyedProcessFunction 加上它的状态、定时器、侧输出三大件。

一句话定位:DataStream 是 Flink 面向”每条记录 + 每个 key 的状态 + 每个时间点的定时器”的命令式底层 API;SQL 够用就别下沉,一旦需要自定义状态、事件时间定时器或多路分流,KeyedProcessFunction 就是你的瑞士军刀。

TL;DR

Table of contents

Open Table of contents

1. 先问一句:SQL 够用吗

在下沉到 DataStream 之前,先立一条纪律:能用 SQL 就别用 DataStream。

Flink SQL / Table API 是声明式的——你说”要什么”(按 campaign 每分钟聚合曝光数),优化器决定”怎么算”(状态怎么存、窗口怎么触发、Join 用什么策略)。它的好处是开发快、可读性高、优化器帮你管状态和窗口,绝大多数实时报表、宽表、聚合都该用 SQL。

那什么时候 SQL 会”卡壳”、必须下沉到 DataStream?看四类信号:

  1. 需要自定义状态的读写时机:不只是”聚合成一个值”,而是要按业务逻辑随时读、随时写、随时清一份 per-key 的状态(如实时用户画像、复杂去重)。
  2. 需要定时器:逻辑依赖”时间流逝”本身——“下单后 15 分钟没支付就告警”、“某 key 静默 5 分钟就清状态”。SQL 里这类”负向事件(某事发生)“极其别扭。
  3. 需要侧输出分流:主结果之外,要把迟到数据、脏数据、命中规则的记录单独送到另一条流去补算/落库/告警。
  4. 复杂事件处理:要精确控制算子对每一条记录、每一个 key 的行为(CEP、状态机)。

这四类的共同点是:你需要对”状态 + 时间 + 分流”的细粒度控制。而这正是 KeyedProcessFunction 的主场。

2. DataStream 的心智:还是 source → transform → sink

好消息是,下沉到 DataStream,架构篇里的心智一点没变:一个作业还是 source → 一串 transformation → sink,数据无界、逐条流动、算子常驻。区别只在于——每个算子的逻辑现在由你用 Java 写死,而不是交给 SQL 优化器生成。

一条有状态 DataStream 作业的骨架图,以及 KeyedProcessFunction 的三大能力。上半部分是横向流水线:KafkaSource(读 orders topic)—fromSource—> keyBy(userId)(按 key 做 hash 重分区)—keyBy—> KeyedProcessFunction(有状态处理,加粗高亮)—process—> Sink(print 或写回 Kafka),呼应架构篇 source→transform→sink 的骨架。下半部分并排三个能力盒,标题为「KeyedProcessFunction 的三大能力:SQL 表达不了、必须下沉 DataStream 的地方」:① ValueState / MapState——按 key 记住中间结果,用于累计、去重、画像;② TimerService——注册 event-time 或 processing-time 定时器,用于超时告警、延迟触发;③ OutputTag 侧输出——在主流之外再开一股流,用于迟到、异常、分流。

DataStream 作业骨架与 SQL 时完全同构,只是算子逻辑改由你手写;真正的价值集中在 KeyedProcessFunction 的状态、定时器、侧输出三大能力上。

一个最小的有状态 DataStream 作业长这样:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
env.enableCheckpointing(5000);              // 开 checkpoint,容错语义与 SQL 一致

env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "orders")
   .map(OrderParser::parse)                 // 解析 JSON
   .filter(o -> o != null)
   .keyBy(o -> o.userId)                    // = SQL 的 GROUP BY userId
   .process(new UserAggFunction())          // 自定义有状态处理
   .print();

env.execute("datastream-order-agg");

注意 keyBy 这一步:它就是 SQL 里的 GROUP BY,按 key 做 hash 重分区,保证同一个 key 的所有记录恒落在同一个并行子任务上——这是后面所有 keyed state 能”按 key 独立记忆”的物理前提。

3. 从 SQL 到 DataStream 的翻译对照

刚从 SQL 过来,可以先建立一张对照表,把熟悉的概念映射到 DataStream:

SQL / Table APIDataStream API说明
GROUP BY keykeyBy(key)按 key hash 重分区,keyed state 的前提
聚合函数 SUM/COUNTValueState + 手动累加 / reduce / aggregate状态从”优化器托管”变成”你显式读写”
TUMBLE / HOP / CUMULATE 窗口.window(...)KeyedProcessFunction + 定时器高层窗口仍可用;要自定义触发就上定时器
Regular / Interval Joinconnect + CoProcessFunction / 双流状态双流关联手动维护两侧状态
无(负向事件很别扭)TimerService + onTimer()”某事没在 X 时间内发生”用定时器天然表达
无(要拆多次 filter)OutputTag 侧输出一次处理分出多股流

一句话:SQL 帮你托管的东西(状态、窗口触发、Join 状态),到了 DataStream 全变成你显式管理——这是代价,也是自由。

4. KeyedProcessFunction:状态 + 定时器 + 侧输出

DataStream 里算子有很多(map / flatMap / window / reduce…),但要论”能干最多脏活”的,是 KeyedProcessFunction<K, I, O>。它必须用在 keyBy 之后,在一个函数里同时给你三样东西:

其它高层算子基本都是它的特例。骨架如下:

public class MyFunction extends KeyedProcessFunction<String, Event, String> {
    private transient ValueState<Long> state;

    @Override
    public void open(Configuration cfg) {
        state = getRuntimeContext().getState(
            new ValueStateDescriptor<>("my-state", Long.class));
    }

    @Override
    public void processElement(Event e, Context ctx, Collector<String> out) throws Exception {
        // 1) 读写当前 key 的状态
        // 2) ctx.timerService().registerEventTimeTimer(ts)  注册定时器
        // 3) ctx.output(LATE_TAG, e)  侧输出
        // 4) out.collect(...)  主流输出
    }

    @Override
    public void onTimer(long ts, OnTimerContext ctx, Collector<String> out) {
        // 定时器到点触发
    }
}

下面用三个能跑的例子,分别把这三大件用起来。

5. ValueState 实战:累计花费与阈值告警

最典型的 keyed state 用法:给每个 key 维护一个累计值,来一条更一条,越过阈值就告警一次。 广告里的”广告主当日花费超 80% 预算提醒”、“用户频控计数”都是这个模式。

public class SpendAlertFunction extends KeyedProcessFunction<String, Order, String> {
    private transient ValueState<Double> totalState;    // 累计花费
    private transient ValueState<Boolean> alertedState; // 是否已告警(去重)

    @Override
    public void open(Configuration cfg) {
        totalState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("total", Double.class));
        alertedState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("alerted", Boolean.class));
    }

    @Override
    public void processElement(Order o, Context ctx, Collector<String> out) throws Exception {
        double total = Optional.ofNullable(totalState.value()).orElse(0.0);
        total += o.amount;
        totalState.update(total);

        out.collect(ctx.getCurrentKey() + " 累计花费=" + total);

        boolean alerted = Optional.ofNullable(alertedState.value()).orElse(false);
        if (total > 1000.0 && !alerted) {           // 越过阈值,且没告警过
            out.collect("⚠️ 告警:" + ctx.getCurrentKey() + " 累计花费突破 1000");
            alertedState.update(true);               // 标记已告警,避免每条都刷
        }
    }
}

两个细节值得记:

⚠️ 按 key 累积的状态若不设 TTL,key 基数一大就会无限膨胀,直接拖垮 checkpoint。这正是状态与 Checkpoint 容错里反复强调的保命开关。

6. TimerService 实战:事件时间支付超时检测

现在上第二件武器——定时器。需求是一个典型的”负向事件”:下单后 15 分钟内没有支付,就发一条超时告警。 注意”没支付”是一件没发生的事,SQL 很难直接表达,而定时器天然合适:下单时注册一个 15 分钟后的定时器,支付到了就把它删掉;没删掉就说明超时了。

事件时间定时器实现"支付超时检测"的时序图。三个参与方从左到右:订单事件流、KeyedProcessFunction、TimerService。主线步骤:1) create 下单事件到达 KeyedProcessFunction;2) 函数把 orderTs 存入 ValueState;3) 函数向 TimerService 注册一个 event-time 定时器(下单时间 +15 分钟)。随后分成两支:绿色分支①(正常)——15 分钟内 pay 支付事件到达,函数调用 deleteEventTimeTimer 取消定时器并清空状态;虚线表示 pay 到达与删除定时器。红色分支②(超时)——Watermark 越过定时器时间仍未支付,TimerService 回调 onTimer(timestamp) 触发函数,函数输出「支付超时」告警(可再经 OutputTag 分流)。核心:定时器由 event-time 的 Watermark 推进触发,而非机器时钟。

“某事没在规定时间内发生”这类负向逻辑,用定时器最自然:下单注册定时器、支付删定时器,到点还没删就是超时。

public class PayTimeoutFunction extends KeyedProcessFunction<String, Event, String> {
    private transient ValueState<Long> orderTsState;  // 下单时间
    private transient ValueState<Long> timerTsState;  // 已注册的定时器时间

    @Override
    public void open(Configuration cfg) {
        orderTsState = getRuntimeContext().getState(new ValueStateDescriptor<>("orderTs", Long.class));
        timerTsState = getRuntimeContext().getState(new ValueStateDescriptor<>("timerTs", Long.class));
    }

    @Override
    public void processElement(Event e, Context ctx, Collector<String> out) throws Exception {
        if ("create".equals(e.type)) {
            orderTsState.update(e.eventTime);
            long timerTs = e.eventTime + 15 * 60 * 1000L;     // 下单 + 15 分钟
            ctx.timerService().registerEventTimeTimer(timerTs); // 注册事件时间定时器
            timerTsState.update(timerTs);
        } else if ("pay".equals(e.type)) {
            Long timerTs = timerTsState.value();
            if (timerTs != null) {
                ctx.timerService().deleteEventTimeTimer(timerTs); // 支付到了:删掉定时器
            }
            orderTsState.clear();
            timerTsState.clear();                                // 正常,清状态
        }
    }

    @Override
    public void onTimer(long ts, OnTimerContext ctx, Collector<String> out) throws Exception {
        // 定时器触发 = 15 分钟内没等到 pay 把它删掉
        if (orderTsState.value() != null) {
            out.collect("⚠️ 订单 " + ctx.getCurrentKey() + " 支付超时(下单后 15 分钟未支付)");
            orderTsState.clear();
            timerTsState.clear();
        }
    }
}

关键点是事件时间定时器由 Watermark 推进触发,而不是机器时钟:只有当 Watermark 越过 timerTs,Flink 才认为”事件时间到了 15 分钟后”并回调 onTimer()。这保证了可重放、可复现——同一批数据重跑一遍,超时判定完全一致,不会因为你今天跑得慢/快而变。这也是为什么定时器一定要配合正确的 Watermark 策略使用。

7. Side Output 实战:迟到 / 异常分流

第三件武器——侧输出(Side Output)。主流只承载”正常结果”,而迟到的、格式错的、命中某条规则的记录,用 OutputTag 分到另一条流,单独补算 / 落库 / 告警。它比”用多个 filter 把流拆好几遍”更高效(只处理一次),也能和窗口的 allowedLateness 迟到数据无缝配合。

public class L15Job {
    // OutputTag 必须是匿名内部类(带上泛型信息),否则类型擦除拿不到
    public static final OutputTag<String> LATE_TAG = new OutputTag<String>("late-orders") {};

    public static void main(String[] args) throws Exception {
        // ... env / source ...
        SingleOutputStreamOperator<String> main = keyed.process(new PayTimeoutFunction());

        main.print();                                  // 主流:正常结果
        main.getSideOutput(LATE_TAG).print();          // 侧流:迟到 / 异常
        env.execute("L15-Timer与SideOutput");
    }
}

// 在 processElement 里判断迟到,走侧输出
if (e.eventTime < ctx.timerService().currentWatermark()) {
    ctx.output(LATE_TAG, "迟到事件: " + e);            // 分流,不进主流
    return;
}

侧输出的价值在于**“不丢、也不污染主流”**:广告计费里,迟到 3 秒的曝光如果直接丢,就是漏计费;直接混进主流又会破坏窗口的准确性。用侧输出把它捞出来单独补算,是”既准确又不丢”的标准做法。

8. 接真实数据源:KafkaSource + ValueState 聚合

把上面这套接到真实数据源——直接消费 Kafka 的订单流,用 DataStream 重算一遍”按用户累计”。这样就和 SQL 版的实时聚合形成对照:同一个 Kafka topic,SQL 和 DataStream 都能算,容错语义(靠 checkpoint)完全一致。

KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
    .setBootstrapServers("kafka:9092")
    .setTopics("orders")
    .setGroupId("datastream-order-agg")
    .setStartingOffsets(OffsetsInitializer.latest())
    .setValueOnlyDeserializer(new SimpleStringSchema())
    .build();

env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "kafka-orders")
   .map(L16::parseOrder).name("parse-json")     // JSON → Order
   .filter(o -> o != null)
   .keyBy(o -> o.userId)
   .process(new UserAggFunction())              // ValueState 累计(同第 5 节模式)
   .print();

env.execute("L16-Kafka订单DataStream聚合");

这里 KafkaSource 默认把 offset 存进 checkpoint(承接端到端 Exactly-Once 的可重放 source)——作业挂掉从 checkpoint 恢复时,offset 和 ValueState 一起回滚、一起重放,累计值不会翻倍。DataStream 并没有牺牲任何容错能力,它和 SQL 共享同一套 checkpoint 机制。

SQL 作业可以在 SQL Client 里跑,DataStream 作业则要编译打包成 jar 再提交。用 Maven 管依赖(flink-streaming-javaflink-connector-kafka 等设为 provided,因为集群 lib 里已有),maven-jar-plugin 指定 mainClassmvn package 打出 jar,再 flink run

# 打包
mvn -q clean package

# 提交到集群(把 target 挂进容器后)
flink run -d /opt/jars/datastream-1.0.jar

# 生产更推荐 Application 模式(main 在 JobManager 上跑,见部署篇)
flink run-application -t kubernetes-application ...

依赖设 provided 是关键:作业 jar 里不该再打包 Flink 运行时,否则和集群 lib 里的版本冲突、jar 还臃肿。提交与部署模式(Session vs Application、Flink on K8s)是部署篇的主题。

10. 动手实战(本地 playground)

在本地 docker Flink 集群上,这三大件都能”边跑边看”:

  1. 起集群docker compose up -d,Web UI localhost:8081
  2. 打包mvn -q clean package,把 target 挂到 JobManager/TaskManager 的 /opt/jars
  3. ValueState 累计告警flink run -d /opt/jars/...L14....jardatagen 造订单,标准输出里能看到累计值一路上涨,突破阈值时打出一条 ⚠️ 告警 后不再刷屏(告警去重生效)。
  4. 支付超时:跑 L15,只发 create 不发 pay,等 Watermark 推过 15 分钟(造数时用较大时间步长加速),Web UI 里能看到 onTimer 触发、主流打出超时告警;补发 pay 的那些 key 则不会告警。
  5. 侧输出:故意发几条时间戳早于当前 Watermark 的迟到事件,观察它们只出现在侧流 print、没污染主流。
  6. 接 Kafka:先跑 SQL 版把订单写进 orders topic,再 flink run L16 的 DataStream 聚合,两者结果对得上——同一份数据、两套 API、一致语义。

观察重点:在 Web UI 里,一个 KeyedProcessFunction 作业的算子链、并行度、checkpoint 大小,和 SQL 作业看起来没有本质区别——因为它们最终都被翻译成同一套算子物理执行(架构篇的三次图变换)。DataStream 只是让你手写了算子逻辑而已。

参考


views
Share this post on:

Previous Post
Flink 深挖 · 状态管理与 Checkpoint 容错
Next Post
Flink 深挖 · 端到端 Exactly-Once 与两阶段提交