Flink 流式处理基础

从流处理定位、时间语义、状态与 Checkpoint 到 Spring Boot 订单实时统计案例,系统梳理 Flink 的核心工作方式

Posted by Ekko on August 20, 2026

这篇笔记的目标,是把 Flink 放回真实实时数据链路里重新理解一遍:它并不只是“另一种大数据框架”,而是一套围绕无界数据流、事件时间、状态管理和故障恢复构建出来的分布式计算系统。

文章重点不放在 API 罗列,而放在一条更稳定的认知主线:为什么很多实时业务最后都会走到 FlinkWatermarkCheckpoint 分别解决了什么问题,状态为什么是流处理系统的核心资产,以及一个 Spring Boot 业务系统如何与 Flink 正确协作,形成可维护的实时统计链路。

参考资料:

官方资料:Apache FlinkStateful Stream ProcessingCheckpointing

连接器资料:Kafka Connector

[TOC]


Flink 可以概括为一套面向有状态流处理的分布式计算引擎。它最核心的能力,不是“跑得快”,而是能够在持续到来的无界数据上,稳定维护计算状态,并在机器故障、进程重启、数据乱序、重复消费这些现实问题存在时,仍然尽量保持结果正确。

如果把这个定义拆开来看,重点有四层:

  • 无界数据流:输入不是一次性读完的批量数据,而是持续不断到来的事件
  • 有状态计算:统计、去重、窗口聚合、模式检测,本质上都依赖状态
  • 事件时间语义:很多业务关心的是“事件发生时间”,不是“程序收到时间”
  • 故障恢复能力:流处理作业可能长期运行数周甚至数月,没有可恢复能力就不具备生产价值

因此,Flink 的强项从来不是“把 SQL 或 Java 代码并行跑起来”这么简单,而是:

在分布式环境里持续处理事件流,同时维护一致的状态视图,并在失败后恢复到可接受的一致性边界。

很多业务系统并不是缺少“计算框架”,而是缺少一套能稳定处理实时变化的系统。典型问题包括:

  • 订单、支付、点击、埋点持续不断到来,不能等到夜间批处理才出结果
  • 同一业务事件可能乱序到达,甚至重复到达
  • 聚合逻辑需要持续维护中间状态,而不是每次全量重算
  • 程序重启后不能把过去几十分钟的统计结果全部丢掉
  • 既希望低延迟,又希望结果具备稳定的一致性

这些问题拼在一起,才是 Flink 真正要解决的对象。

Flink、Kafka、Hive、Spring Batch 最容易混淆的地方

Flink 常常和消息队列、离线数仓、批处理框架一起出现,因此边界必须先讲清楚。

系统 主要职责 擅长场景 不擅长场景
Kafka 事件传输与缓冲 解耦生产者和消费者、承接高吞吐消息流 复杂状态计算、窗口聚合、规则计算
Flink 实时流处理与状态计算 实时聚合、去重、告警、CEP、维表关联 纯消息存储、强事务 OLTP
Hive 离线数仓与批量分析 历史数据扫描、宽表加工、离线统计 秒级实时计算
Spring Batch 应用内定时批处理 定时清洗、批量导入导出、单体/中小规模任务 高吞吐无界数据流处理

一句话概括就是:

  • Kafka 负责把事件可靠地送到下游
  • Flink 负责把事件持续算出结果
  • Hive 负责在离线数仓里做历史分析
  • Spring Batch 更偏应用内定时批任务,而不是分布式实时流计算

理解边界比背定义更重要。Flink 并不天然适合下面这些问题:

  • 只在每天凌晨跑一次、结果晚几小时也没关系的统计任务
  • 高并发事务写入、行级锁、强一致更新这类 OLTP 场景
  • 简单到用数据库 GROUP BY 就能完成的轻量聚合
  • 没有稳定消息源、没有重放能力、没有持久化状态存储的“伪实时”脚本

如果问题本质上是“数据持续到来,结果要持续更新,而且失败后还能恢复”,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 负责真正执行算子

这里的 DAGDirected Acyclic Graph 的缩写,中文通常翻译为“有向无环图”。

在初始理解阶段,可以先将它视为一张“数据处理流程图”:

  • 有向:数据有明确流向,例如从 Kafka 进入,再经过清洗、聚合,最后写入 MySQL
  • 无环:主处理链路不会绕一圈又回到原点,否则同一条数据就可能在图里无限打转
  • :不是只做一步,而是由多个处理节点按顺序连起来

如果把一条订单支付消息的处理过程写成最朴素的流程,大致就是:

Kafka 读取订单消息 -> 解析 JSON -> 过滤支付事件 -> 按店铺分组 -> 做分钟窗口聚合 -> 写入 MySQL

这整条流程,就是一个非常典型的 DAGFlink 作业在逻辑上就是这样一张图:图上的每个节点都是一个算子,边表示数据流向。这里的“算子”可以先理解为“处理数据的一步操作”。常见节点包括:

  • 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 时间过长
  • 同步调用外部服务,把流式作业拖成串行阻塞程序

事件时间和处理时间的差别,不是术语差别,而是结果差别

实时业务里最容易被低估的,是“到底按什么时间算”。

时间语义 含义 优点 常见问题
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 偏移量、广播配置、算子级快照

其中最常见的,是挂在 KeyedProcessFunctionRichFlatMapFunction、窗口聚合上的 ValueStateMapStateListState

Checkpoint 不是“备份”,而是分布式一致性快照

在初始理解阶段,可以先把 Checkpoint 理解成“定期存档”。但这里保存的不是单个变量,而是整条作业链路在某个时刻的一致状态。

没有状态快照,流处理作业一旦失败,就只能从头再来。问题在于,很多实时场景根本不允许从头再来,因为:

  • 上游保留时间有限,不一定能完整重放
  • 状态规模很大,从头回算成本极高
  • 对外已经输出过结果,重复计算可能造成重复写入

FlinkCheckpoint 机制,本质上是在整个数据流图上做一致性快照。可以用下图理解:

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 对应的状态和消费位置。只要上游可重放、下游语义匹配,就能尽量逼近无故障执行时的结果。

很多介绍会简单写成“Flink 支持 exactly-once”,这句话本身不假,但如果不讲条件,很容易造成误解。

端到端结果能否做到严格一次,取决于整个链路:

环节 要求
Source 能按 checkpoint 恢复到正确消费位置,支持重放
Flink 内部状态 checkpoint 成功持久化,可在故障后恢复
Sink 能做到事务提交、两阶段提交,或至少业务层幂等

因此,很多工程实践会选择一套更容易维护的折中方案:

  • Flink 内部依赖 Checkpoint
  • 输出端使用幂等 upsert
  • 业务表通过唯一键覆盖更新

这不一定是最严格的理论 exactly-once,但往往是更现实、也更稳定的生产方案。

四、常见坑点与失效场景

把处理时间当事件时间,结果会悄悄漂移

如果直接按处理时间开窗口,网络延迟、消息堆积、消费抖动都会改变窗口归属。最开始结果可能看起来“差得不多”,但在活动流量、跨机房链路或 Kafka 积压场景下,偏差会迅速放大。

只开窗口,不定义迟到策略,晚到数据会直接丢业务价值

窗口关闭之后,晚到事件怎么处理,必须在设计阶段明确:

  • 允许一定迟到时间并继续修正窗口
  • 放入侧输出流,交给补偿链路
  • 明确丢弃,但要在业务上能解释

如果这一步没定义,线上通常会出现“统计数总比业务库少一点,但查不出哪里错”的情况。

去重只靠数据库唯一键,实时链路会被写库拖垮

很多团队在早期实现阶段,会把重复数据问题完全留给数据库唯一键处理。这种方式短期内可以运行,但长期通常会遇到两个问题:

  • 大量重复写请求把数据库写放大
  • Flink 作业无法在流内提前拦截重复事件,资源浪费明显

更合理的方式是:流内先用状态做轻量去重,落库侧再用唯一键做最后一道幂等保护。

key 选得不好,状态再正确也跑不快

错误示例包括:

  • 按固定常量 keyBy
  • 按极度倾斜字段 keyBy
  • 按变化过于频繁的组合字段做 key,导致状态爆炸

keyBy 选型本质上是在平衡三件事:

  • 业务语义是否成立
  • 数据是否均匀
  • 状态规模是否可控

Checkpoint 配了,但存储不可靠,恢复仍然是空谈

开发环境把 checkpoint 放在本地目录往往没问题,但生产上如果 JobManagerTaskManager 或宿主机丢失,本地文件也会一起丢掉。真正可靠的做法,通常是把 checkpoint 放到持久化共享存储里,例如 HDFS 或对象存储。

从 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,而不是发送时机器当前时间

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

版本号应按当前集群版本和连接器兼容矩阵统一管理,不建议手工混搭。

下面这段代码展示最核心的处理骨架:读取 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 的原因
实时大盘、分钟级指标统计 连续到来的事件需要持续聚合
风控规则、异常检测、CEP 需要状态、时序模式、低延迟
实时 ETL、日志清洗、事件路由 需要高吞吐、可扩展的数据流转换
流批一体数据管道 同一套引擎承接实时与部分批式处理
CDC 后的实时宽表或实时数仓链路 需要维护动态状态并实时下游分发
场景 更合适的方案
每天跑一次的报表任务 HiveSpark、调度批任务
低频数据同步、小批量导入导出 Spring Batch、定时脚本
纯消息转发,没有状态计算 Kafka Connect、轻量消费者
简单实时查询,数据库本身足够承担 直接数据库聚合或缓存方案

是否引入 Flink,更值得优先回答的是下面这几个问题:

  • 数据是不是持续不断到来
  • 结果是不是要持续更新
  • 计算是不是离不开状态
  • 故障恢复是不是业务刚需
  • 乱序、重复、晚到是不是现实存在的问题

如果这几个问题大多都回答“是”,那么 Flink 往往就不是锦上添花,而是系统能力的基础设施。