实时计算 Flink 版

Realtime Compute for Apache Flink

数据一产生就立刻处理,让报表和告警跟上当下,而不是等到第二天才看到结果。

存储与数据出现时间 · 2016 年前后按作业占用的计算资源与运行时长计费,作业不开就不花钱。流式计算实时双十一大屏Flink
看官方文档

01 它是干嘛的

Flink 版是阿里云托管的流式计算服务,基于开源 Apache Flink。它处理的是「源源不断、永不停止」的数据流:数据一到就计算,几秒内出结果。与批量计算「攒够了再算」相对,它适合监控、实时大屏、实时风控等对延迟敏感的场景。

02 为什么会有它

从「天级报表」到「秒级大屏」:批处理为什么不够用

传统的数据统计是「批处理」:每天凌晨把前一天的数据一次性算完,早上你看到的是昨天的报表。这像每天收一次信——不是不能用,但任何需要「现在就要知道」的场景都赶不上:库存告急、系统异常、订单欺诈,等你第二天看到,损失已经发生。

双十一这类场景把矛盾推到极致。成千上万人同时下单,运营需要一块实时大屏,显示此刻的成交额、订单量、地域分布,延迟只能以秒计。如果用批处理,数据要攒一段时间再算,大屏就是「过去的画面」,失去了指挥价值。

流式计算正是为「数据一到就算」而生:数据像水流一样不断涌入,处理程序持续接收、持续计算、持续输出。它不需要等数据攒齐,因此延迟可以做到秒级甚至毫秒级。Flink 是目前最主流的流式计算引擎之一,而阿里把它大规模用于双十一的实时链路,并把这套能力做成托管服务对外提供。

03 它怎么工作

流式计算的最小模型是三段:Source(数据源)负责把数据接进来,Operator(算子)负责在流动中做计算,Sink(数据汇)负责把结果送出去。数据一直在管道里流动,不停留、不攒批。

流式计算:数据一到就算,不等攒批① 持续流入② 算子链③ 落地④ 直接展示消息队列 / 埋点每秒几万条事件解析与清洗过滤脏数据、补字段窗口聚合每 5 秒算一次结果表 / 告警写回数据库触发告警实时大屏秒级刷新批处理是「攒一天再算」,流计算是「来一条算一条」。双十一大屏、风控拦截靠的都是流计算。持续不断的数据流计算核心
顺着箭头看「Source → Operator → Sink」这条流动的线:数据从不停留,这是流式计算与批处理最本质的区别。

一条数据从产生到出现在大屏上:Source 到 Operator 再到 Sink

  1. 1
    ① Source 接入数据

    数据来源可以是消息队列、日志、数据库变更流等,Source 负责持续不断地把新数据读进来。

  2. 2
    ② 数据在算子间流动

    算子是处理逻辑,如过滤、映射、聚合、窗口统计;数据像流水线一样从一个算子流向下一个。

  3. 3
    ③ 按时间窗口统计

    流是无限的,要统计「每分钟成交额」就用窗口把无穷的流切成一段段有限的数据来算。

  4. 4
    ④ 状态与容错

    聚合的中间结果保存在状态里;系统会定期做检查点,故障后从最近一次检查点恢复,不丢数据。

  5. 5
    ⑤ Sink 输出结果

    结果写进数据库、消息队列或直接推给大屏,秒级更新,供人实时查看与决策。

04 谁在用它

实时大屏与监控

大促、直播、赛事等活动现场,实时展示成交、在线人数、地域分布等关键指标。

实时风控与反欺诈

交易产生的瞬间就判断是否异常,疑似盗刷当场拦截,而不是事后追查。

实时告警

对系统指标、业务指标做持续计算,一旦越界立即触发告警通知。

实时数据同步

把数据库的变更实时同步到下游系统,让各系统数据保持一致。

实时特征计算

为推荐和风控模型持续计算最新特征,让模型的判断基于「此刻」而非昨天的数据。

05 怎么用

会写 SQL 就能使用它的主流模式(Flink SQL),真正的难点在于理解流与窗口的概念。

  1. 01开通实时计算 Flink 版并创建一个工作空间与项目。
  2. 02配置数据源(如消息队列 Kafka)和数据目标(如数据库、大屏)。
  3. 03用 Flink SQL 或 Java 程序编写处理逻辑,定义窗口与聚合方式。
  4. 04提交作业并观察运行状态,按数据量调整并发度与资源。
  5. 05配置告警与自动重启策略,保证作业长期稳定运行。
用 Flink SQL 统计每分钟成交额sql
-- 从订单流中,按每分钟统计成交金额,结果写入下游表
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 关键概念

流式计算
数据一到就处理、不攒批,追求低延迟,适合实时场景。
批处理
先把数据攒起来,再一次性计算,吞吐高但结果有延迟。
窗口
把无限的流按时间切成一段段有限数据来统计,如每分钟、每五分钟。
检查点
定期保存计算状态,作业故障后从这里恢复,保证数据不丢。

07 容易混淆的对比

Spark Streaming同样处理流数据,但常按「微批」方式工作,延迟通常高于真正逐条处理的 Flink。
自建 Flink 集群完全自主可控,但要自己负责版本、扩容与故障处理,运维成本高。
图软件图鉴

纯静态站点,无后端、无追踪。内容为中文原创撰写,用于帮助非技术读者理解主流软件产品。

关于本站

全部内容在构建时生成
产品名称与商标归各自公司所有

© 2026 软件图鉴Nuxt 4 · Vue 3 · Tailwind 4 · 纯静态