前面几篇里,无论是时间窗口驱动的实时大盘,还是端到端 Exactly-Once 的实时宽表,能用 Flink SQL 表达的我们都用 SQL——它声明式、上手快,状态与窗口还有优化器替你管着。但广告实时链路一走到深水区,总会撞上一类 SQL 写不出来的需求:「某个广告主 15 分钟内没完成转化就告警」「这条曝光迟到了 3 秒,别丢,单独送去补算」「按用户维护一份随时能读写的实时画像」,要的都是对状态、时间与分流的细粒度控制,而这恰恰是声明式 SQL 给不了的。
这时候就该下沉到 Flink 的底层 API——DataStream。本篇承接前几篇的 SQL 与真实数据源,把 DataStream 最核心的一套武器讲透、跑通:KeyedProcessFunction 以及它自带的状态、定时器、侧输出三大件。
一句话定位:DataStream 是 Flink 面向「每条记录 + 每个 key 的状态 + 每个时间点的定时器」的命令式底层 API;SQL 够用就别下沉,一旦需要自定义状态、事件时间定时器或多路分流,
KeyedProcessFunction就是那把瑞士军刀。
TL;DR
- 先问一句:SQL 够用吗? 能用 Flink SQL / Table API 表达的聚合、窗口与 Join,优先交给 SQL——声明式、状态由优化器托管、维护成本最低。DataStream 是 SQL 写不出来时才下沉的底层 API,不是默认选项。
- 必须下沉 DataStream 的四类信号:① 需要自定义状态的读写时机,而不只是聚合出一个值;② 需要事件时间或处理时间定时器来做超时与延迟触发;③ 需要侧输出把迟到、异常、命中规则的记录从主流里单独分出来;④ 复杂事件处理,要精确控制算子对每条记录的行为。
- DataStream 心智没变:还是
source → 一串 transformation → sink(承接架构篇),区别只在于每个算子的逻辑改由你用 Java 写死,而不是交给 SQL 优化器生成。 keyBy= SQL 的GROUP BY:按 key 做 hash 重分区,保证同 key 恒落同一并行子任务,这是一切 keyed state 的物理前提。- 两类状态:Keyed State 的作用域是每个 key,只能在
keyBy之后使用;Operator State 的作用域是每个并行实例。Keyed 原语有ValueState / ListState / MapState / ReducingState / AggregatingState(详见状态与 Checkpoint)。 KeyedProcessFunction是最强算子:一个函数里同时拿到 state(按 key 记住中间结果)+TimerService(注册定时器)+Context(侧输出与时间戳),其它高层算子(map、window)都可以看作它的特例。ValueState做累计与告警:把每个 key 的累计值存进ValueState,来一条更一条、越过阈值发一次告警,实时频控与预算超额提醒都是这个模式。TimerService做超时检测:注册 event-time 定时器,由 Watermark 推进触发onTimer();「下单后 15 分钟未支付就告警」这类基于时间流逝的逻辑,SQL 很别扭,定时器很自然。OutputTag侧输出做分流:主流之外再开一股或多股流,专收迟到事件、脏数据、命中某条规则的记录——比用多次filter反复拆流高效,也能和窗口的allowedLateness配合兜住迟到数据。- 接真实数据源:
KafkaSource+keyBy+KeyedProcessFunction就能把真实数据源那套 Kafka 订单流用 DataStream 重算一遍,状态一致性依旧靠 checkpoint 兜底,容错语义与 SQL 完全一致。 - 打包提交:DataStream 作业用 Maven 打成 jar 后由
flink run提交,生产环境走 Application 模式(见部署篇)。
Table of contents
Open Table of contents
- 1. 先问一句:SQL 够用吗
- 2. DataStream 的心智:还是 source → transform → sink
- 3. 从 SQL 到 DataStream 的翻译对照
- 4. KeyedProcessFunction:状态 + 定时器 + 侧输出
- 5. ValueState 实战:累计花费与阈值告警
- 6. TimerService 实战:事件时间支付超时检测
- 7. Side Output 实战:迟到 / 异常分流
- 8. 接真实数据源:KafkaSource + ValueState 聚合
- 9. 打包与提交:Maven + flink run
- 10. 动手实战(本地 playground)
- 参考
1. 先问一句:SQL 够用吗
在下沉到 DataStream 之前,先立一条纪律:能用 SQL 就别用 DataStream。
Flink SQL / Table API 是声明式的——你只说清楚「要什么」(按 campaign 每分钟聚合曝光数),至于状态怎么存、窗口何时触发、Join 走哪种策略,全部交给优化器决定。开发快、可读性高、状态与窗口有人替你兜底,绝大多数实时报表、实时宽表与聚合本就该停在这一层。
那么 SQL 会在什么地方卡壳、非下沉不可?归纳起来是四类信号:
- 需要自定义状态的读写时机:不只是把一串数值聚合成一个结果,而是要按业务逻辑随时读、随时写、随时清一份 per-key 的状态,实时用户画像与复杂去重都属此类。
- 需要定时器:逻辑依赖「时间流逝」本身——下单后 15 分钟没支付就告警、某个 key 静默 5 分钟就清状态。这类「负向事件(某件事没发生)」用 SQL 表达极其别扭。
- 需要侧输出分流:主结果之外,还要把迟到数据、脏数据、命中规则的记录单独送去补算、落库或告警。
- 复杂事件处理:要精确控制算子对每一条记录、每一个 key 的行为,比如 CEP 与状态机。
四类信号指向的是同一件事:你需要对「状态 + 时间 + 分流」的细粒度控制。而这正是 KeyedProcessFunction 的主场,本篇后面的篇幅基本都围绕它展开。
2. DataStream 的心智:还是 source → transform → sink
好消息是,下沉到 DataStream 之后,架构篇里那套心智一点没变:一个作业仍然是 source → 一串 transformation → sink,数据无界、逐条流动、算子常驻集群。真正的区别只有一处——每个算子的逻辑现在由你用 Java 写死,而不是交给 SQL 优化器生成。
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 独立记忆」——这不是约定,而是物理前提,没有 keyBy 就没有 keyed state。
3. 从 SQL 到 DataStream 的翻译对照
刚从 SQL 过来,最快的入门方式是先建立一张对照表,把熟悉的概念逐条映射到 DataStream:
| SQL / Table API | DataStream API | 说明 |
|---|---|---|
GROUP BY key | keyBy(key) | 按 key hash 重分区,keyed state 的前提 |
聚合函数 SUM/COUNT | ValueState + 手动累加 / reduce / aggregate | 状态从优化器托管变成你显式读写 |
TUMBLE / HOP / CUMULATE 窗口 | .window(...) 或 KeyedProcessFunction + 定时器 | 高层窗口仍可用,要自定义触发才上定时器 |
| Regular / Interval Join | connect + CoProcessFunction / 双流状态 | 双流关联需手动维护两侧状态 |
| 无(负向事件很别扭) | TimerService + onTimer() | 某事没在 X 时间内发生,用定时器天然表达 |
| 无(要拆多次 filter) | OutputTag 侧输出 | 一次处理分出多股流 |
一句话概括这张表:SQL 帮你托管的东西(状态、窗口触发、Join 状态),到了 DataStream 全部变成你显式管理的对象——这既是成本,也是自由,而下一节的 KeyedProcessFunction 就是行使这份自由的入口。
4. KeyedProcessFunction:状态 + 定时器 + 侧输出
DataStream 的算子不少(map / flatMap / window / reduce……),但要论谁能干最多脏活,非 KeyedProcessFunction<K, I, O> 莫属。它必须接在 keyBy 之后,并在同一个函数里一次性交给你三样东西:
- 状态(State):在
open()里初始化ValueState / MapState / ListState等描述符,运行时由 Flink 按当前记录的 key 自动切换状态视图,你写的永远是「当前这个 key 的那一份」。 TimerService:通过ctx.timerService()注册 event-time 或 processing-time 定时器,到点回调onTimer()。Context:拿到当前记录的 key 与时间戳,以及侧输出入口ctx.output(tag, value)。
正因为这三者同时在手,其它高层算子基本都可以看作它的特例。骨架如下:
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 到 §7 就沿着这个骨架层层加码:先只用状态,再给状态配上定时器,最后把定时器判出来的异常记录经侧输出分走。
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); // 标记已告警,避免每条都刷
}
}
}
两个细节值得记住:
- 状态是 per-key 的:
totalState.value()拿到的永远是当前这条记录所属 key 的那一份,Flink 会随 key 自动切换状态视图,你不必也不该自己维护一个Map<key, total>——手搓的 Map 既进不了 checkpoint,也没法在扩缩并行度时随 key group 迁移。 - 告警去重也交给状态:用
alertedState记住这个 key 已经告过警,阈值突破后就不会每来一条都刷屏。这种「状态里再存一个开关」的写法,在告警与频控场景里非常常见。
⚠️ 按 key 累积的状态若不设 TTL,key 基数一大就会无限膨胀,直接拖垮 checkpoint。这正是状态与 Checkpoint 容错里反复强调的保命开关。
6. TimerService 实战:事件时间支付超时检测
光有状态还不够:上一节的告警靠的是「来了一条记录」这个触发点,可业务里有大量逻辑恰恰是没有记录到来时才该触发。于是第二件武器登场:定时器。需求是一个典型的负向事件:下单后 15 分钟内没有支付就发超时告警。 既然「没支付」是一件没发生的事,SQL 无从下手,定时器却天然合适——下单时注册一个 15 分钟后的定时器,支付到了就把它删掉,到点还没被删掉,就说明超时了。
「某事没在规定时间内发生」这类负向逻辑,用定时器最自然:下单注册定时器、支付删定时器,到点还没删就是超时。
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 策略上,上一篇里那些乱序容忍度与空闲分区(withIdleness)的取舍,在这里会原样兑现成告警的早报或漏报。
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;
}
注意这里复用的正是 §6 的 PayTimeoutFunction:同一个算子,主流继续吐超时告警,侧流承接迟到事件。侧输出的价值就在于既不丢、也不污染主流——广告计费里,迟到 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 机制,差别只在算子逻辑由谁书写。
9. 打包与提交:Maven + flink run
SQL 作业可以在 SQL Client 里直接跑,DataStream 作业则要先编译打包成 jar 再提交。依赖用 Maven 管理,flink-streaming-java、flink-connector-kafka 等一律设为 provided(集群 lib 里已经有了),maven-jar-plugin 指定 mainClass,mvn 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 与 Application 怎么选、Flink on K8s 怎么落地,是部署篇的主题。
10. 动手实战(本地 playground)
在本地 docker Flink 集群上,这三大件都能边跑边看:
- 起集群:
docker compose up -d,Web UIlocalhost:8081。 - 打包:
mvn -q clean package,把target挂到 JobManager/TaskManager 的/opt/jars。 - ValueState 累计告警:
flink run -d /opt/jars/...L14....jar,datagen造订单,标准输出里能看到累计值一路上涨,突破阈值时打出一条⚠️ 告警后不再刷屏,说明告警去重生效。 - 支付超时:跑 L15,只发
create不发pay,等 Watermark 推过 15 分钟(造数时用较大时间步长加速),Web UI 里能看到onTimer触发、主流打出超时告警;补发过pay的那些 key 则安静无声。 - 侧输出:故意发几条时间戳早于当前 Watermark 的迟到事件,观察它们只出现在侧流 print 里,没有污染主流。
- 接 Kafka:先跑 SQL 版把订单写进
orderstopic,再flink runL16 的 DataStream 聚合,两者结果应当对得上——同一份数据、两套 API、一致语义。
观察重点:在 Web UI 里,一个
KeyedProcessFunction作业的算子链、并行度与 checkpoint 大小,和 SQL 作业看起来没有本质区别——因为两者最终都被翻译成同一套算子物理执行(架构篇的三次图变换)。DataStream 只是让你手写了算子逻辑而已。
参考
- Apache Flink. DataStream API Overview:DataStream 编程模型、算子与作业结构官方总览。
- Apache Flink. Process Function (Low-level Operations):
ProcessFunction / KeyedProcessFunction、TimerService与状态的权威说明。 - Apache Flink. Working with State:Keyed/Operator State 原语、
StateTtlConfig与状态后端。 - Apache Flink. Side Outputs:
OutputTag侧输出的用法与迟到数据处理。 - Apache Flink. Kafka Connector(KafkaSource / KafkaSink):DataStream 侧 Kafka 连接器与 offset / checkpoint 语义。