服务器与大带宽专家 · 持牌IDC/CDN/ISP服务商
简米科技官网JIANMI TECH
资讯 2026-09-16 更新于 2026-09-16 简米科技 3,826 字 9 分钟阅读

流批一体如何统一处理实时与历史数据?流批一体架构思路

导读用同一套SQL或代码逻辑,同时跑通实时流数据和离线历史数据,从而避免维护两套计算引擎和两套业务口径,为什么流批一体思路越来越被数据团队接受传统Lambda架构的痛点正在倒逼变革过去处理实时与历史数据,最常见的做法是搭两套链路,实时链路:Flink、Storm或Kafka Streams,数据从消息队列进,秒级出……

用同一套SQL或代码逻辑,同时跑通实时流数据和离线历史数据,从而避免维护两套计算引擎和两套业务口径。

为什么流批一体思路越来越被数据团队接受

传统Lambda架构的痛点正在倒逼变革

过去处理实时与历史数据,最常见的做法是搭两套链路。

  • 实时链路:Flink、Storm或Kafka Streams,数据从消息队列进,秒级出结果。
  • 离线链路:Hive、Spark SQL或MapReduce,凌晨跑批,T+1产出报表。
  • 两套代码:业务逻辑写两遍,实时一个口径,离线一个口径。
  • 数据对不上:实时大屏显示100万,离线报表出来98万,业务方追着数仓问原因。

这种架构维护成本高,数据修正时,实时链路要补数,离线链路要重跑,开发人员不敢动口径,因为一动就要改两个地方。

流批一体的本质是把批处理看作有界流

行业内普遍把流处理定义为“无界数据”,批处理定义为“有界数据”。

  • 实时流是一条永不停止的数据河流,事件不断到达。
  • 历史数据是一段已经固定下来的区间,可以全部加载。
  • 流批一体就是让同一条SQL或同一套算子,既能跑在无界流上,也能跑在有界数据上。

用Flink的视角看,批处理只是流处理的一个特例,同一份逻辑,切换数据源类型即可。

流批一体实时数仓怎么选型才能少走弯路

选型先看三件事

第一,团队技术栈。 如果团队以SQL为主,优先选Flink SQL或Spark SQL,如果团队有较强的Java/Scala能力,可以考虑DataStream API或Dataset API做更细粒度控制。

第二,数据规模和延迟要求。 实时大屏要求秒级延迟,实时数仓多要求分钟级,历史数据回溯频率是每天一次还是每周一次,直接影响选型。

第三,开源生态和运维成本。 开源自建灵活但需要专人维护,云托管省心但按量付费,华东地区不少数据团队更倾向云上托管,因为可以省去机器采购和集群调优时间。

Flink流批一体和Spark Structured Streaming对比

流批一体如何统一处理实时与历史数据?流批一体架构思路

维度

Flink Spark Structured Streaming
流处理原生性 原生流设计,低延迟 微批次为主,延迟稍高
状态管理 成熟的键控状态和定时器 有状态算子但生态较弱
时间语义 事件时间处理完善 事件时间支持但配置略复杂
批处理能力 批为有界流,统一良好 批处理是强项,历史久
SQL成熟度 流批SQL语法统一 SQL主要用于批,流SQL在跟上
生态集成 与Kafka、Pulsar、Iceberg等集成成熟 与Delta Lake、Hive集成更好

业内专家指出,流批一体选型没有绝对答案,关键看实时链路是否长期占据主导地位,如果实时需求强,Flink是更自然的选择,如果离线数仓历史包袱重,Spark Structured Streaming过渡成本更低。

用同一套逻辑处理实时与历史数据的落地路径

第一步:统一数据源和格式

实时数据进Kafka或Pulsar,历史数据落到Apache Iceberg、Hudi或Delta Lake,两边定义相同的表Schema,字段名、类型、约束保持一致。

  • 订单表实时源:Kafka topic order_realtime,JSON格式。
  • 订单表历史源:Iceberg表 ods.order_history,Parquet格式。
  • 统一注册到Flink Catalog,使用同一张逻辑表名 order_source

第二步:用SQL抽象业务逻辑

以下是一段简化的Flink SQL示例,展示实时与历史共用聚合逻辑。

-- 实时表
CREATE TABLE order_realtime (
  order_id STRING,
  amount DECIMAL(10,2),
  order_time TIMESTAMP(3),
  WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH ('connector'='kafka', ...);
-- 历史表
CREATE TABLE order_history (
  order_id STRING,
  amount DECIMAL(10,2),
  order_time TIMESTAMP(3)
) WITH ('connector'='iceberg', ...);
-- 同一套聚合逻辑,通过视图切换
CREATE VIEW unified_orders AS
SELECT order_id, amount, order_time FROM order_realtime
UNION ALL
SELECT order_id, amount, order_time FROM order_history;

流批一体如何统一处理实时与历史数据?流批一体架构思路

然后对 unified_orders 执行相同的金额求和或去重逻辑。

第三步:处理状态与时间差异

实时流依赖watermark推进事件时间,历史数据批处理时,没有持续的watermark。

操作路径如下:

  • 实时模式:设置watermark生成策略,允许迟到数据。
  • 批处理模式:使用处理时间或让watermark等于最大事件时间,避免等待。
  • 在Flink SQL中,可以通过作业参数 execution.runtime-mode 切换 STREAMINGBATCH

第四步:结果输出保持一致

实时结果写入OLAP引擎如StarRocks或ClickHouse,历史修正结果也要写入同一张结果表,这样才能保证大屏和报表数字一致。

  • 实时链路每5分钟触发一次聚合,结果upsert到结果表。
  • 历史回溯每天凌晨跑一次同口径聚合,覆盖当天部分或全部数据。
  • 用主键去重和版本号控制,避免重复累计。

流批一体在电商实时大屏场景下的应用实例

某电商大促期间,大屏需要展示实时成交额、订单量、对比去年同期,传统做法实时链路算一次,离线链路算去年同期,两边数字经常对不上。

流批一体的落地步骤:

  1. 实时订单写入Kafka,历史订单同步到Iceberg。
  2. 定义统一的订单聚合SQL,按分钟或小时粒度计算成交额。
  3. 实时部分读取Kafka,历史部分读取Iceberg中去年同期的分区。
  4. 同一套SQL对两个数据集分别执行聚合。
  5. 合并结果写入大屏展示表,实时部分秒级更新,历史部分固定不变。

这样大屏上的实时数字和离线对比数字来自同一份业务逻辑,口径天然一致。

流批一体数据仓库价格贵不贵?成本拆开看

开源自建成本

  • 机器成本:需要Flink集群、Kafka集群、数据湖存储,数据量中等规模的话,三到五台节点可以起步。
  • 人力成本:至少需要一名熟悉Flink和SQL的工程师,负责集群维护和作业开发。
  • 软件成本:开源组件免费,但版本升级、漏洞修复需要自己跟进。
  • 流批一体如何统一处理实时与历史数据?流批一体架构思路

云上托管成本

多数云厂商按计算资源CU或作业数计费,实时作业和批处理作业分开计费,流批一体可以减少作业数量,因为同一套逻辑不用部署两遍。

  • 实时作业按流量或CU小时计费。
  • 批处理作业按扫描数据量或运行时长计费。
  • 如果历史回溯频率高,云上成本可能比自建更可控。

行业共识认为,流批一体在数据量达到一定规模后,比维护两套Lambda架构更省钱,因为开发和运维成本下降明显。

流批一体的常见误区

一套代码完全不用改

实时和批处理仍然存在微小差异,比如watermark、状态清理策略、数据迟到处理,但整体逻辑可以复用绝大部分。

所有场景都适合流批一体

低频小数据量的报表需求,传统批处理可能更简单,流批一体更适合实时性要求高且历史回溯频繁的场景。

流批一体等于实时数仓

流批一体是处理架构思路,实时数仓是数据组织方式,两者经常搭配,但不完全等同。

流批一体的本质不是技术炫技,而是让业务口径只定义一次,当实时与历史数据用同一套逻辑处理,数据打架的问题自然消失。

Q&A

流批一体和Lambda架构的主要区别是什么?

Lambda架构需要同时维护实时层和批处理层两套代码,两套逻辑独立演进,容易产生口径不一致,流批一体用统一的API或SQL处理有界和无界数据,开发和维护成本更低。

流批一体实时数仓适合小公司吗?

适合中等数据量以上、有明确实时需求且团队具备一定SQL能力的公司,小公司数据量小、实时需求弱时,简单的批处理或单条流处理可能更省事,不能为了流批一体而强行上复杂架构。

流批一体用Flink还是Spark更合适?

多数情况下Flink在流批一体上更成熟,因为Flink原生按流设计,批处理作为有界流特例,状态管理和事件时间更完善,Spark Structured Streaming从批处理起家,微批次模式会导致实时延迟略高,不过近年来也在持续改进流处理能力,选择时优先看团队现有技术栈和实时链路占比。

分享本文
本文为 简米科技官网 原创,已由运维技术专家审核。转载请注明来源:原文链接
售前咨询 服务热线 售后 邮箱