前面几篇里,无论是时间窗口还是端到端 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(记住东西)+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。
那什么时候 SQL 会”卡壳”、必须下沉到 DataStream?看四类信号:
- 需要自定义状态的读写时机:不只是”聚合成一个值”,而是要按业务逻辑随时读、随时写、随时清一份 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 独立记忆”的物理前提。
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 全变成你显式管理——这是代价,也是自由。
4. KeyedProcessFunction:状态 + 定时器 + 侧输出
DataStream 里算子有很多(map / flatMap / window / reduce…),但要论”能干最多脏活”的,是 KeyedProcessFunction<K, I, O>。它必须用在 keyBy 之后,在一个函数里同时给你三样东西:
- 状态(State):在
open()里初始化ValueState / MapState / ListState等,按当前 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. 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>(那样既不容错也不能 rescale)。 - 告警去重也用状态:用
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 策略使用。
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 机制。
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 vs 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 语义。