这篇笔记的目标,是把
Flink放回真实实时数据链路里重新理解一遍:它并不只是“另一种大数据框架”,而是一套围绕无界数据流、事件时间、状态管理和故障恢复构建出来的分布式计算系统。
文章重点不放在 API 罗列,而放在一条更稳定的认知主线:为什么很多实时业务最后都会走到
Flink,Watermark和Checkpoint分别解决了什么问题,状态为什么是流处理系统的核心资产,以及一个Spring Boot业务系统如何与Flink正确协作,形成可维护的实时统计链路。
参考资料:
官方资料:Apache Flink 、 Stateful Stream Processing 、 Checkpointing
连接器资料:Kafka Connector
[TOC]
一、先回答 Flink 到底是什么
Flink 可以概括为一套面向有状态流处理的分布式计算引擎。它最核心的能力,不是“跑得快”,而是能够在持续到来的无界数据上,稳定维护计算状态,并在机器故障、进程重启、数据乱序、重复消费这些现实问题存在时,仍然尽量保持结果正确。
如果把这个定义拆开来看,重点有四层:
无界数据流:输入不是一次性读完的批量数据,而是持续不断到来的事件有状态计算:统计、去重、窗口聚合、模式检测,本质上都依赖状态事件时间语义:很多业务关心的是“事件发生时间”,不是“程序收到时间”故障恢复能力:流处理作业可能长期运行数周甚至数月,没有可恢复能力就不具备生产价值
因此,Flink 的强项从来不是“把 SQL 或 Java 代码并行跑起来”这么简单,而是:
在分布式环境里持续处理事件流,同时维护一致的状态视图,并在失败后恢复到可接受的一致性边界。
Flink 主要解决什么问题
很多业务系统并不是缺少“计算框架”,而是缺少一套能稳定处理实时变化的系统。典型问题包括:
- 订单、支付、点击、埋点持续不断到来,不能等到夜间批处理才出结果
- 同一业务事件可能乱序到达,甚至重复到达
- 聚合逻辑需要持续维护中间状态,而不是每次全量重算
- 程序重启后不能把过去几十分钟的统计结果全部丢掉
- 既希望低延迟,又希望结果具备稳定的一致性
这些问题拼在一起,才是 Flink 真正要解决的对象。
Flink、Kafka、Hive、Spring Batch 最容易混淆的地方
Flink 常常和消息队列、离线数仓、批处理框架一起出现,因此边界必须先讲清楚。
| 系统 | 主要职责 | 擅长场景 | 不擅长场景 |
|---|---|---|---|
Kafka |
事件传输与缓冲 | 解耦生产者和消费者、承接高吞吐消息流 | 复杂状态计算、窗口聚合、规则计算 |
Flink |
实时流处理与状态计算 | 实时聚合、去重、告警、CEP、维表关联 | 纯消息存储、强事务 OLTP |
Hive |
离线数仓与批量分析 | 历史数据扫描、宽表加工、离线统计 | 秒级实时计算 |
Spring Batch |
应用内定时批处理 | 定时清洗、批量导入导出、单体/中小规模任务 | 高吞吐无界数据流处理 |
一句话概括就是:
Kafka负责把事件可靠地送到下游Flink负责把事件持续算出结果Hive负责在离线数仓里做历史分析Spring Batch更偏应用内定时批任务,而不是分布式实时流计算
Flink 不是什么
理解边界比背定义更重要。Flink 并不天然适合下面这些问题:
- 只在每天凌晨跑一次、结果晚几小时也没关系的统计任务
- 高并发事务写入、行级锁、强一致更新这类 OLTP 场景
- 简单到用数据库
GROUP BY就能完成的轻量聚合 - 没有稳定消息源、没有重放能力、没有持久化状态存储的“伪实时”脚本
如果问题本质上是“数据持续到来,结果要持续更新,而且失败后还能恢复”,Flink 很合适;如果问题本质上是“低频批任务”,就没有必要额外引入实时计算引擎。
二、Flink 的运行机制为什么适合实时计算
先看一个最小运行视图
graph LR
A[Spring Boot 业务系统]
B[Kafka Topic]
C[Source]
D[KeyBy]
E[Window / Process]
F[State]
G[Sink]
H[MySQL / Redis]
I[JobManager]
J[TaskManager]
A --> B
B --> C
C --> D
D --> E
E --> G
E --> F
G --> H
I --> J
J --> C
J --> D
J --> E
J --> G
这张图里有两条主线:
- 业务数据主线:事件从上游进入
Kafka,由Flink读取、分区、计算、输出 - 运行控制主线:
JobManager负责调度作业,TaskManager负责真正执行算子
Flink 作业本质上是一个持续运行的 DAG
这里的 DAG 是 Directed Acyclic Graph 的缩写,中文通常翻译为“有向无环图”。
在初始理解阶段,可以先将它视为一张“数据处理流程图”:
有向:数据有明确流向,例如从Kafka进入,再经过清洗、聚合,最后写入MySQL无环:主处理链路不会绕一圈又回到原点,否则同一条数据就可能在图里无限打转图:不是只做一步,而是由多个处理节点按顺序连起来
如果把一条订单支付消息的处理过程写成最朴素的流程,大致就是:
Kafka 读取订单消息 -> 解析 JSON -> 过滤支付事件 -> 按店铺分组 -> 做分钟窗口聚合 -> 写入 MySQL
这整条流程,就是一个非常典型的 DAG。Flink 作业在逻辑上就是这样一张图:图上的每个节点都是一个算子,边表示数据流向。这里的“算子”可以先理解为“处理数据的一步操作”。常见节点包括:
Source:数据进入作业的输入端Map / FlatMap / Filter:单条记录变换KeyBy:按照某个字段把同类数据分到一起,方便后面按这个字段做统计Window / ProcessFunction:在一段时间范围内,或者带着中间状态去处理数据Sink:结果离开作业并写入外部系统的输出端
如果只保留最核心的职责划分,可以概括为:
Source:把数据读进来KeyBy:把同一类数据放到一起Window:按时间段做统计Sink:把结果写出去
这套模型的重要点在于:它不是“执行一次就结束”的 DAG,而是持续运行、持续接收数据、持续维护状态的 DAG。
并行度、分区、KeyBy 决定了数据怎么被拆开计算
在 Flink 里,几乎所有扩展性问题最后都会落到并行度和分区模型上。初始理解阶段,可以先把它归结为一个更基础的问题:
一条持续运行的数据流水线,应当由多少个并行执行单元共同处理,以及数据应按什么规则分配到这些执行单元上。
| 概念 | 作用 | 需要特别关注什么 |
|---|---|---|
parallelism |
决定一个算子开多少个并行实例,可以先理解成“这一处理步骤有多少个并行执行单元” | 不是越大越好,受数据量、资源、分区数限制 |
slot |
TaskManager 上资源隔离单元,可以先理解成“机器上预留出来的执行位置” |
影响任务如何被调度到机器上 |
partition |
上游数据如何分片,也就是消息本身被切成多少份并行流 | 直接影响消费并发与热点 |
keyBy |
按 key 重分区,让同 key 数据进入同一实例,本质上是在说“同类数据必须交给同一个地方统计” | 选错 key 会导致状态错乱或数据倾斜 |
例如订单实时统计按 shopId 聚合时,必须先 keyBy(shopId),否则同一家店铺的数据会分散在不同并行实例上,中间状态无法正确累加。
Operator Chain 为什么能提升吞吐
很多人第一次看 Flink Web UI 时会发现,一个看起来写了很多 API 的作业,实际执行图未必有那么多节点。这通常与 Operator Chain 有关。
在初始理解阶段,可以先把 Operator Chain 理解成“把若干轻量处理步骤串接为一段连续执行链路”。Flink 会尽量把能够串起来的轻量算子放进同一个线程链路里执行,减少线程切换、序列化和网络传输开销。对于简单的 map -> filter -> flatMap 这类连续变换,这个优化非常有效。
但这也意味着一件事:
流处理系统的性能瓶颈,不只来自业务逻辑本身,也来自数据在算子之间如何流动。
Backpressure 才是实时链路最真实的压力信号
只要下游消费速度跟不上上游生产速度,背压就会沿着算子链逐级向上游传播。其直接表现通常是:下游越慢,上游受到的阻塞也越明显。
一个典型链路的背压传播可以表示为:
graph LR
A[Kafka Source]
B[解析与过滤]
C[聚合算子]
D[JDBC Sink]
D -->|写库变慢| C
C -->|缓冲堆积| B
B -->|处理受阻| A
这也是为什么很多实时作业表面看是“Flink 慢”,实际根因却在:
- 下游数据库写入能力不足
- 热 key 导致某个并行实例成为单点瓶颈
- 状态过大导致 checkpoint 时间过长
- 同步调用外部服务,把流式作业拖成串行阻塞程序
三、时间语义、状态与一致性为什么是 Flink 的核心
事件时间和处理时间的差别,不是术语差别,而是结果差别
实时业务里最容易被低估的,是“到底按什么时间算”。
| 时间语义 | 含义 | 优点 | 常见问题 |
|---|---|---|---|
Processing Time |
事件被算子处理时所在机器的当前时间 | 配置简单、延迟低 | 乱序和网络抖动会直接影响结果 |
Event Time |
事件在业务中真实发生的时间 | 结果更贴近业务事实 | 需要处理乱序、延迟、水位线 |
对于订单、支付、日志、埋点这类场景,生产上更可靠的选择通常是 Event Time。原因很直接:业务关心的是支付发生在 10:01:02,不是这条消息因为网络抖动在 10:01:11 才被消费者拿到。
Watermark 解决的是“什么时候可以认为一段时间基本收齐了”
无界流没有天然结束点,所以窗口什么时候触发,必须靠额外机制来判断。Watermark 的本质,是流处理系统对“时间进度”的一个声明。
可以把它理解成:
当前系统大致认为,时间戳早于某个点的数据已经基本到齐,可以开始对相关窗口做计算。
一个简化过程如下:
sequenceDiagram
participant S as Source
participant O as Operator
participant W as Window
S->>O: event ts=10:00:03
S->>O: event ts=10:00:08
S->>O: watermark=10:00:10
O->>W: 时间推进到 10:00:10
W->>W: 关闭已到期窗口
S->>O: late event ts=10:00:05
O->>W: 判定是否允许迟到或侧输出
这里最关键的不是记住 API,而是形成三个判断:
- 水位线不是“当前最新事件时间”,而是“系统认为已经足够安全推进的时间”
- 窗口不是等到现实时间过去才关,而是等到
watermark越过窗口边界才关 - 晚到数据不是没有,而是必须提前定义如何处理
状态为什么是流处理的真正资产
很多流处理逻辑看起来像是在“一条条处理数据”,但真正决定结果是否正确的,是这些数据背后的状态。这里的“状态”可以先理解为“程序为了继续计算而必须保留的中间结果”。
例如下面这些操作,本质上都离不开状态:
- 按用户统计最近 10 分钟点击数
- 对同一订单的重复支付事件做去重
- 按设备 ID 维护最近一次上报时间
- 做窗口聚合,累计每个 key 的中间结果
Flink 常见状态可以粗分为两类:
| 状态类型 | 适用对象 | 典型用途 |
|---|---|---|
Keyed State |
已经过 keyBy 的数据流,可以理解成“每个 key 维护各自独立的状态” |
去重、累计、会话状态、每个 key 的中间值 |
Operator State |
算子实例级别状态,可以理解成“算子实例自身维护的运行信息” | Source 偏移量、广播配置、算子级快照 |
其中最常见的,是挂在 KeyedProcessFunction、RichFlatMapFunction、窗口聚合上的 ValueState、MapState、ListState。
Checkpoint 不是“备份”,而是分布式一致性快照
在初始理解阶段,可以先把 Checkpoint 理解成“定期存档”。但这里保存的不是单个变量,而是整条作业链路在某个时刻的一致状态。
没有状态快照,流处理作业一旦失败,就只能从头再来。问题在于,很多实时场景根本不允许从头再来,因为:
- 上游保留时间有限,不一定能完整重放
- 状态规模很大,从头回算成本极高
- 对外已经输出过结果,重复计算可能造成重复写入
Flink 的 Checkpoint 机制,本质上是在整个数据流图上做一致性快照。可以用下图理解:
graph LR
A[Kafka Partition 0]
B[Kafka Partition 1]
C[Barrier]
D[算子状态快照]
E[Checkpoint Storage]
A --> C
B --> C
C --> D
D --> E
这个过程真正保存的,不只是“某个对象当前值”,而是:
- 各个有状态算子的当前状态
- 数据源已经处理到什么位置
- 整张作业图在某个逻辑时刻的一致性切面
当任务失败后,Flink 会把作业恢复到最近一次成功 Checkpoint 对应的状态和消费位置。只要上游可重放、下游语义匹配,就能尽量逼近无故障执行时的结果。
Exactly-Once 不是只靠 Flink 一家就能保证
很多介绍会简单写成“Flink 支持 exactly-once”,这句话本身不假,但如果不讲条件,很容易造成误解。
端到端结果能否做到严格一次,取决于整个链路:
| 环节 | 要求 |
|---|---|
| Source | 能按 checkpoint 恢复到正确消费位置,支持重放 |
| Flink 内部状态 | checkpoint 成功持久化,可在故障后恢复 |
| Sink | 能做到事务提交、两阶段提交,或至少业务层幂等 |
因此,很多工程实践会选择一套更容易维护的折中方案:
Flink内部依赖Checkpoint- 输出端使用幂等
upsert - 业务表通过唯一键覆盖更新
这不一定是最严格的理论 exactly-once,但往往是更现实、也更稳定的生产方案。
四、常见坑点与失效场景
把处理时间当事件时间,结果会悄悄漂移
如果直接按处理时间开窗口,网络延迟、消息堆积、消费抖动都会改变窗口归属。最开始结果可能看起来“差得不多”,但在活动流量、跨机房链路或 Kafka 积压场景下,偏差会迅速放大。
只开窗口,不定义迟到策略,晚到数据会直接丢业务价值
窗口关闭之后,晚到事件怎么处理,必须在设计阶段明确:
- 允许一定迟到时间并继续修正窗口
- 放入侧输出流,交给补偿链路
- 明确丢弃,但要在业务上能解释
如果这一步没定义,线上通常会出现“统计数总比业务库少一点,但查不出哪里错”的情况。
去重只靠数据库唯一键,实时链路会被写库拖垮
很多团队在早期实现阶段,会把重复数据问题完全留给数据库唯一键处理。这种方式短期内可以运行,但长期通常会遇到两个问题:
- 大量重复写请求把数据库写放大
- Flink 作业无法在流内提前拦截重复事件,资源浪费明显
更合理的方式是:流内先用状态做轻量去重,落库侧再用唯一键做最后一道幂等保护。
key 选得不好,状态再正确也跑不快
错误示例包括:
- 按固定常量
keyBy - 按极度倾斜字段
keyBy - 按变化过于频繁的组合字段做 key,导致状态爆炸
keyBy 选型本质上是在平衡三件事:
- 业务语义是否成立
- 数据是否均匀
- 状态规模是否可控
Checkpoint 配了,但存储不可靠,恢复仍然是空谈
开发环境把 checkpoint 放在本地目录往往没问题,但生产上如果 JobManager、TaskManager 或宿主机丢失,本地文件也会一起丢掉。真正可靠的做法,通常是把 checkpoint 放到持久化共享存储里,例如 HDFS 或对象存储。
将 Flink 作业直接嵌入 Spring Boot 进程,往往会带来生命周期错位
从 Java 语言层面看,Flink 作业可以放在某个应用模块里编译和打包;但从运行方式看,它和常规 Spring Boot Web 应用属于两类系统:
Spring Boot偏同步接口、应用容器、请求响应生命周期Flink偏长期运行作业、资源调度、状态恢复、savepoint 升级
因此更合理的工程边界通常是:
Spring Boot负责业务系统、接口、消息生产、结果查询Flink作为独立作业部署和运维
五、实战案例:Spring Boot 订单系统的实时 GMV 统计
场景定义
假设有一个电商订单系统,希望在运营后台展示:
- 每个店铺最近一分钟支付订单数
- 每个店铺最近一分钟支付金额
- 数据延迟控制在数秒级
- 支持事件乱序、重复投递、作业失败后恢复
如果直接在 Spring Boot 应用里做同步统计,会遇到几个问题:
- 接口线程不适合维护长生命周期窗口状态
- 支付回调可能重试,事件会重复
- 订单事件可能乱序到达
- 应用重启后内存状态全部丢失
因此更适合拆成一条标准实时链路:
graph LR
A[Spring Boot 订单服务]
B[Kafka: order-events]
C[Flink 实时作业]
D[MySQL: ads_shop_gmv_1m]
E[运营看板]
A --> B
B --> C
C --> D
D --> E
事件模型设计
实时系统要先把事件建模清楚,否则后面所有时间语义和去重逻辑都会漂。
订单事件可以定义为:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
public class OrderEvent {
private String eventId;
private String orderId;
private Long shopId;
private String eventType;
private java.math.BigDecimal amount;
private long eventTime;
public String getEventId() {
return eventId;
}
public void setEventId(String eventId) {
this.eventId = eventId;
}
public String getOrderId() {
return orderId;
}
public void setOrderId(String orderId) {
this.orderId = orderId;
}
public Long getShopId() {
return shopId;
}
public void setShopId(Long shopId) {
this.shopId = shopId;
}
public String getEventType() {
return eventType;
}
public void setEventType(String eventType) {
this.eventType = eventType;
}
public java.math.BigDecimal getAmount() {
return amount;
}
public void setAmount(java.math.BigDecimal amount) {
this.amount = amount;
}
public long getEventTime() {
return eventTime;
}
public void setEventTime(long eventTime) {
this.eventTime = eventTime;
}
}
这里有两个字段尤其关键:
eventId:用于去重eventTime:用于事件时间窗口计算
Spring Boot 订单服务如何发送事件
实际工程里,Spring Boot 更适合作为业务事件生产者,而不是 Flink 运行容器。下面给一个简化的发送示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
@Service
public class OrderEventPublisher {
private final KafkaTemplate<String, OrderEvent> kafkaTemplate;
public OrderEventPublisher(KafkaTemplate<String, OrderEvent> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public void publishPaidEvent(String orderId, Long shopId, BigDecimal amount, Instant paidAt) {
OrderEvent event = new OrderEvent();
event.setEventId(UUID.randomUUID().toString());
event.setOrderId(orderId);
event.setShopId(shopId);
event.setEventType("PAID");
event.setAmount(amount);
event.setEventTime(paidAt.toEpochMilli());
kafkaTemplate.send("order-events", String.valueOf(shopId), event);
}
}
这段代码刻意保留了两个工程细节:
- Kafka key 直接使用
shopId,有利于同店铺事件落到相对稳定的分区 - 事件时间来自业务时间
paidAt,而不是发送时机器当前时间
Flink 作业的基本依赖
如果使用 Java DataStream API,常见依赖大致如下:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
<dependencies>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>${flink.kafka.connector.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-jdbc</artifactId>
<version>${flink.version}</version>
</dependency>
</dependencies>
版本号应按当前集群版本和连接器兼容矩阵统一管理,不建议手工混搭。
Flink 作业主链路
下面这段代码展示最核心的处理骨架:读取 Kafka、定义 watermark、过滤支付事件、按店铺聚合、落 MySQL。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
public class ShopMinuteGmvJob {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(4);
env.enableCheckpointing(30000L);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(10000L);
env.getCheckpointConfig().setCheckpointTimeout(60000L);
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
KafkaSource<String> source = KafkaSource.<String>builder()
.setBootstrapServers("localhost:9092")
.setTopics("order-events")
.setGroupId("shop-minute-gmv")
.setStartingOffsets(OffsetsInitializer.latest())
.setValueOnlyDeserializer(new SimpleStringSchema())
.build();
WatermarkStrategy<OrderEvent> watermarkStrategy =
WatermarkStrategy.<OrderEvent>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, timestamp) -> event.getEventTime());
DataStream<OrderEvent> orderEventStream = env
.fromSource(source, WatermarkStrategy.noWatermarks(), "kafka-order-source")
.map(Jsons::readOrderEvent)
.assignTimestampsAndWatermarks(watermarkStrategy)
.filter(event -> "PAID".equals(event.getEventType()));
DataStream<OrderEvent> deduplicatedStream = orderEventStream
.keyBy(OrderEvent::getEventId)
.process(new EventIdDeduplicateFunction());
DataStream<ShopMinuteGmv> resultStream = deduplicatedStream
.keyBy(OrderEvent::getShopId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.allowedLateness(Time.seconds(30))
.aggregate(new ShopMinuteGmvAggregate(), new ShopMinuteGmvWindowFunction());
resultStream.addSink(MySqlUpsertSink.build());
env.execute("shop-minute-gmv-job");
}
}
这条链路里最值得注意的不是代码量,而是每一步的职责边界:
| 步骤 | 作用 | 为什么不能省 |
|---|---|---|
map(Jsons::readOrderEvent) |
反序列化业务事件 | 没有稳定事件模型,后续无法统一处理 |
assignTimestampsAndWatermarks |
定义事件时间与乱序容忍度 | 不定义时间语义,窗口结果会漂移 |
keyBy(eventId) 去重 |
拦截重复消息 | 支付回调重试和消息重复消费都很常见 |
keyBy(shopId) 聚合 |
汇总店铺级指标 | 同店铺数据必须进入同一逻辑分区 |
allowedLateness |
接住一部分晚到数据 | 否则窗口关闭后迟到事件没有去处 |
upsert sink |
把结果幂等写出 | 故障恢复和窗口修正都可能带来重复输出 |
用状态实现轻量去重
很多实时去重并不需要上来就做全局 Bloom Filter,先把单事件幂等做好,收益就已经很大。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
public class EventIdDeduplicateFunction
extends KeyedProcessFunction<String, OrderEvent, OrderEvent> {
private transient ValueState<Boolean> seenState;
@Override
public void open(Configuration parameters) {
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(org.apache.flink.api.common.time.Time.hours(2))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.neverReturnExpired()
.build();
ValueStateDescriptor<Boolean> descriptor =
new ValueStateDescriptor<>("seen-event", Types.BOOLEAN);
descriptor.enableTimeToLive(ttlConfig);
seenState = getRuntimeContext().getState(descriptor);
}
@Override
public void processElement(
OrderEvent value,
Context ctx,
Collector<OrderEvent> out) throws Exception {
if (Boolean.TRUE.equals(seenState.value())) {
return;
}
seenState.update(Boolean.TRUE);
out.collect(value);
}
}
这个实现有三个工程含义:
- 以
eventId作为 key,避免同一事件重复进入下游窗口 - 给状态加
TTL,避免去重状态无限膨胀 - 去重发生在窗口之前,减少无效聚合和无效写库
窗口聚合对象
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
public class ShopMinuteGmv {
private Long shopId;
private Instant windowStart;
private Instant windowEnd;
private long paidOrderCount;
private BigDecimal paidAmount;
public void addAmount(BigDecimal delta) {
if (paidAmount == null) {
paidAmount = BigDecimal.ZERO;
}
paidAmount = paidAmount.add(delta);
}
// getter / setter 省略
}
对应的聚合函数可以写成:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
public class ShopMinuteGmvAggregate
implements AggregateFunction<OrderEvent, ShopMinuteGmv, ShopMinuteGmv> {
@Override
public ShopMinuteGmv createAccumulator() {
ShopMinuteGmv acc = new ShopMinuteGmv();
acc.setPaidAmount(BigDecimal.ZERO);
return acc;
}
@Override
public ShopMinuteGmv add(OrderEvent value, ShopMinuteGmv acc) {
acc.setShopId(value.getShopId());
acc.setPaidOrderCount(acc.getPaidOrderCount() + 1);
acc.addAmount(value.getAmount());
return acc;
}
@Override
public ShopMinuteGmv getResult(ShopMinuteGmv acc) {
return acc;
}
@Override
public ShopMinuteGmv merge(ShopMinuteGmv a, ShopMinuteGmv b) {
a.setPaidOrderCount(a.getPaidOrderCount() + b.getPaidOrderCount());
a.addAmount(b.getPaidAmount());
return a;
}
}
窗口补充函数负责补齐窗口起止时间:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
public class ShopMinuteGmvWindowFunction
extends ProcessWindowFunction<ShopMinuteGmv, ShopMinuteGmv, Long, TimeWindow> {
@Override
public void process(
Long shopId,
Context context,
Iterable<ShopMinuteGmv> elements,
Collector<ShopMinuteGmv> out) {
ShopMinuteGmv result = elements.iterator().next();
result.setShopId(shopId);
result.setWindowStart(Instant.ofEpochMilli(context.window().getStart()));
result.setWindowEnd(Instant.ofEpochMilli(context.window().getEnd()));
out.collect(result);
}
}
MySQL 表设计与幂等落库
如果结果用于运营看板,一个足够实用的设计是使用窗口主键做幂等覆盖:
1
2
3
4
5
6
7
8
9
CREATE TABLE ads_shop_gmv_1m (
window_start DATETIME NOT NULL,
window_end DATETIME NOT NULL,
shop_id BIGINT NOT NULL,
paid_order_count BIGINT NOT NULL,
paid_amount DECIMAL(18, 2) NOT NULL,
updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (window_start, window_end, shop_id)
);
对应 sink 可以采用 upsert:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
public final class MySqlUpsertSink {
private MySqlUpsertSink() {
}
public static SinkFunction<ShopMinuteGmv> build() {
return JdbcSink.sink(
"""
INSERT INTO ads_shop_gmv_1m
(window_start, window_end, shop_id, paid_order_count, paid_amount, updated_at)
VALUES (?, ?, ?, ?, ?, NOW())
ON DUPLICATE KEY UPDATE
paid_order_count = VALUES(paid_order_count),
paid_amount = VALUES(paid_amount),
updated_at = NOW()
""",
(ps, value) -> {
ps.setTimestamp(1, Timestamp.from(value.getWindowStart()));
ps.setTimestamp(2, Timestamp.from(value.getWindowEnd()));
ps.setLong(3, value.getShopId());
ps.setLong(4, value.getPaidOrderCount());
ps.setBigDecimal(5, value.getPaidAmount());
},
JdbcExecutionOptions.builder()
.withBatchSize(200)
.withBatchIntervalMs(1000)
.withMaxRetries(3)
.build(),
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl("jdbc:mysql://localhost:3306/realtime")
.withDriverName("com.mysql.cj.jdbc.Driver")
.withUsername("root")
.withPassword("root")
.build()
);
}
}
这里特意使用 ON DUPLICATE KEY UPDATE,就是为了让同一个窗口的重复输出变成覆盖更新,而不是追加脏数据。
这个案例真正解决了哪些生产问题
把整个案例串起来,实际上解决了五类常见问题:
| 生产问题 | 对应设计 |
|---|---|
| 支付回调重复发送 | eventId 去重状态 |
| 事件存在 5 秒内乱序 | Watermark + Event Time |
| 少量晚到数据需要修正窗口 | allowedLateness(30s) |
| 作业故障后需要恢复 | Checkpoint |
| 落库结果可能重复发送 | 主键 upsert 幂等写入 |
这也是 Flink 真正适合实时统计的原因:不是因为代码看起来高级,而是因为它把这些现实问题纳入了同一条计算链路。
六、边界分析:什么时候该用 Flink,什么时候不必上 Flink
适合使用 Flink 的场景
| 场景 | 选择 Flink 的原因 |
|---|---|
| 实时大盘、分钟级指标统计 | 连续到来的事件需要持续聚合 |
| 风控规则、异常检测、CEP | 需要状态、时序模式、低延迟 |
| 实时 ETL、日志清洗、事件路由 | 需要高吞吐、可扩展的数据流转换 |
| 流批一体数据管道 | 同一套引擎承接实时与部分批式处理 |
| CDC 后的实时宽表或实时数仓链路 | 需要维护动态状态并实时下游分发 |
不一定适合使用 Flink 的场景
| 场景 | 更合适的方案 |
|---|---|
| 每天跑一次的报表任务 | Hive、Spark、调度批任务 |
| 低频数据同步、小批量导入导出 | Spring Batch、定时脚本 |
| 纯消息转发,没有状态计算 | Kafka Connect、轻量消费者 |
| 简单实时查询,数据库本身足够承担 | 直接数据库聚合或缓存方案 |
是否引入 Flink,更值得优先回答的是下面这几个问题:
- 数据是不是持续不断到来
- 结果是不是要持续更新
- 计算是不是离不开状态
- 故障恢复是不是业务刚需
- 乱序、重复、晚到是不是现实存在的问题
如果这几个问题大多都回答“是”,那么 Flink 往往就不是锦上添花,而是系统能力的基础设施。