01 它是干嘛的
Flink 版是阿里云托管的流式计算服务,基于开源 Apache Flink。它处理的是「源源不断、永不停止」的数据流:数据一到就计算,几秒内出结果。与批量计算「攒够了再算」相对,它适合监控、实时大屏、实时风控等对延迟敏感的场景。
02 为什么会有它
从「天级报表」到「秒级大屏」:批处理为什么不够用
传统的数据统计是「批处理」:每天凌晨把前一天的数据一次性算完,早上你看到的是昨天的报表。这像每天收一次信——不是不能用,但任何需要「现在就要知道」的场景都赶不上:库存告急、系统异常、订单欺诈,等你第二天看到,损失已经发生。
双十一这类场景把矛盾推到极致。成千上万人同时下单,运营需要一块实时大屏,显示此刻的成交额、订单量、地域分布,延迟只能以秒计。如果用批处理,数据要攒一段时间再算,大屏就是「过去的画面」,失去了指挥价值。
流式计算正是为「数据一到就算」而生:数据像水流一样不断涌入,处理程序持续接收、持续计算、持续输出。它不需要等数据攒齐,因此延迟可以做到秒级甚至毫秒级。Flink 是目前最主流的流式计算引擎之一,而阿里把它大规模用于双十一的实时链路,并把这套能力做成托管服务对外提供。
03 它怎么工作
流式计算的最小模型是三段:Source(数据源)负责把数据接进来,Operator(算子)负责在流动中做计算,Sink(数据汇)负责把结果送出去。数据一直在管道里流动,不停留、不攒批。
一条数据从产生到出现在大屏上:Source 到 Operator 再到 Sink
- 1① Source 接入数据
数据来源可以是消息队列、日志、数据库变更流等,Source 负责持续不断地把新数据读进来。
- 2② 数据在算子间流动
算子是处理逻辑,如过滤、映射、聚合、窗口统计;数据像流水线一样从一个算子流向下一个。
- 3③ 按时间窗口统计
流是无限的,要统计「每分钟成交额」就用窗口把无穷的流切成一段段有限的数据来算。
- 4④ 状态与容错
聚合的中间结果保存在状态里;系统会定期做检查点,故障后从最近一次检查点恢复,不丢数据。
- 5⑤ Sink 输出结果
结果写进数据库、消息队列或直接推给大屏,秒级更新,供人实时查看与决策。
04 谁在用它
大促、直播、赛事等活动现场,实时展示成交、在线人数、地域分布等关键指标。
交易产生的瞬间就判断是否异常,疑似盗刷当场拦截,而不是事后追查。
对系统指标、业务指标做持续计算,一旦越界立即触发告警通知。
把数据库的变更实时同步到下游系统,让各系统数据保持一致。
为推荐和风控模型持续计算最新特征,让模型的判断基于「此刻」而非昨天的数据。
05 怎么用
会写 SQL 就能使用它的主流模式(Flink SQL),真正的难点在于理解流与窗口的概念。
- 01开通实时计算 Flink 版并创建一个工作空间与项目。
- 02配置数据源(如消息队列 Kafka)和数据目标(如数据库、大屏)。
- 03用 Flink SQL 或 Java 程序编写处理逻辑,定义窗口与聚合方式。
- 04提交作业并观察运行状态,按数据量调整并发度与资源。
- 05配置告警与自动重启策略,保证作业长期稳定运行。
-- 从订单流中,按每分钟统计成交金额,结果写入下游表
INSERT INTO dws_gmv_per_minute
SELECT
TUMBLE_START(order_time, INTERVAL '1' MINUTE) AS window_start,
COUNT(*) AS order_cnt,
SUM(amount) AS gmv
FROM ods_order_stream
GROUP BY TUMBLE(order_time, INTERVAL '1' MINUTE);避坑提示
- !窗口大小与延迟要一起权衡:窗口越长结果越稳但越不实时,越短越实时但抖动越明显。
- !流式作业是长期运行的,要为它配置监控告警与自动恢复,否则挂掉无人知晓。
- !数据乱序(迟到)很常见,要用水位线等机制处理,否则统计会漏掉迟到的数据。
06 关键概念
- 流式计算
- 数据一到就处理、不攒批,追求低延迟,适合实时场景。
- 批处理
- 先把数据攒起来,再一次性计算,吞吐高但结果有延迟。
- 窗口
- 把无限的流按时间切成一段段有限数据来统计,如每分钟、每五分钟。
- 检查点
- 定期保存计算状态,作业故障后从这里恢复,保证数据不丢。