批流一体并不是简单地把批处理任务和流处理任务塞进同一个集群,而是让一套计算引擎、一套存储格式、一套元数据体系同时服务两种计算模式。要理解它的价值,得先看清传统架构的问题:Lambda架构下,批处理链路用Hive或Spark,实时链路用Flink或Storm,同一个业务指标要在两条链路上各实现一次,代码逻辑重复,口径漂移的风险随时存在。每当业务规则变更,两条链路必须同步修改,稍有疏忽就会出现白天报表和实时大屏数字对不上的尴尬局面。

批流一体的核心架构演进:从Lambda到Kappa
Lambda架构是早期大数据平台的主流选择,它把数据链路分成批处理层、速度层和服务层。批处理层定期全量计算保证准确性,速度层处理增量数据保证时效性,服务层负责把两边的计算结果合并对外提供查询。这种架构的致命弱点在于逻辑双写:同一个统计口径要在批和流两套代码里实现,Java写的Flink作业和SQL写的Hive任务很难保证完全等价,尤其是涉及窗口聚合、去重、迟到数据处理这类复杂逻辑时,两边结果出现细微差异几乎是常态。
Kappa架构提出了另一种思路:只保留流处理这一条链路,需要历史重算时把消息队列的消息重放一遍。这在理论上消除了双写问题,但实际落地时对消息队列的容量和回溯能力要求极高,Kafka的 retention 配置要开到足够长,存储成本直线上升,而且一旦消息被清理或者格式升级时没有做好兼容,历史数据就彻底无法重放。
现代批流一体架构可以看作两者的融合升级。以Flink为例,它本身就支持流模式和批模式两种执行方式,同一份DataStream或SQL代码,设置不同的 execution.runtime-mode 即可切换。配合数据湖存储的 snapshot 机制,实时链路写入增量数据,批处理链路按快照读取全量数据,两侧共享同一份物理存储和同一套表结构定义,这才是真正意义上的统一。
// 同一份Flink SQL同时支持批和流两种模式
String sql = "INSERT INTO dws_order_summary " +
"SELECT date_format(order_time, 'yyyy-MM-dd') AS dt, " +
" region, count(*) AS order_cnt, sum(amount) AS gmv " +
"FROM ods_order " +
"GROUP BY date_format(order_time, 'yyyy-MM-dd'), region";
TableEnvironment tEnv = TableEnvironment.create(
EnvironmentSettings.newInstance().build());
// 流模式:持续消费Kafka增量数据
tEnv.getConfiguration().setString("execution.runtime-mode", "streaming");
// 批模式:按快照读取数据湖全量数据做历史回补
// tEnv.getConfiguration().setString("execution.runtime-mode", "batch");
tEnv.executeSql(sql);这种模式下,代码维护成本直接减半,更关键的是批和流读取的是同一份底层数据文件,数据口径天然一致,不需要再花费大量人力去做两条链路的结果比对。
存储层统一:数据湖表格式是批流一体的基石
光有计算引擎统一还不够,存储层如果仍然是流写Kafka、批读写Hive表,那数据还是要落地两份。数据湖表格式(Table Format)解决的就是这个问题,目前主流方案有Iceberg、Hudi和Delta Lake三家。它们的核心能力是给对象存储或HDFS上的裸文件加上一层表格式的抽象,提供ACID事务、schema演进、快照隔离和时间旅行等能力。
Iceberg在这方面的设计比较受Flink社区青睐。它把表的元数据分成三层:metadata file记录表的全局信息,manifest list管理快照,manifest file管理数据文件。每次Flink以流模式写入时,checkpoint提交一次就会生成一个新的snapshot,批处理作业读取时指定snapshot id就能拿到一个事务一致的视图,读写互不阻塞。这种快照机制正是批流共存的物理基础。
-- 创建Iceberg表,批流共用同一张表
CREATE CATALOG iceberg_catalog WITH (
'type'='iceberg',
'catalog-type'='hive',
'uri'='thrift://hive-metastore:9083'
);
CREATE TABLE iceberg_catalog.dws.ods_order (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(16, 2),
order_time TIMESTAMP(3),
region STRING
) PARTITIONED BY (days(order_time))
WITH (
'format-version'='2',
'write.upsert.enabled'='true'
);Hudi的优势在于原生支持upsert语义,对变更频繁的业务表(比如订单状态流转)比较友好,Copy on Write和Merge on Read两种表类型可以按查询性能和写入延迟灵活权衡。Delta Lake则与Spark生态绑定较深,如果团队主力引擎是Spark,选Delta的整合成本最低。选型时不必纠结哪个绝对更好,更重要的是评估团队现有技术栈和社区支持力度。
元数据管理:批流一体集群的治理中枢
元数据是批流一体架构里最容易被轻视、却最影响长期可维护性的部分。传统Hive Metastore作为事实标准存在明显局限:它对表级别的元数据支持没问题,但对文件级别的元数据、快照信息、流式写入的commit记录管理能力很弱,而且高并发场景下的锁竞争一直是痛点,Flink高频checkpoint提交元数据时容易碰到性能瓶颈。
Iceberg这类表格式自带的元数据层实际上分担了Hive Metastore的职责。Iceberg的元数据文件本身就记录了完整的schema历史、分区演进和快照链,Hive Metastore退化成一个目录索引的角色,只负责表的发现和权限挂载。这种分层设计让元数据的读写不再全部压在Metastore一个点上,集群规模扩大后依然能保持稳定。
除了技术层面的元数据,治理层面的元数据管理同样重要,主要包含三块内容。第一是数据血缘:作业之间的依赖关系要能自动采集,解析Flink SQL的输入输出表,构建从ODS到DWS的血缘图谱,一旦上游表结构变更,能快速评估影响范围。第二是数据字典:业务口径、字段含义、负责人信息要和表结构一起管理,避免时间久了没人说得清某个指标的具体含义。第三是权限与审计:统一元数据中心对接权限体系,无论是批作业还是流作业访问表,都走同一套鉴权逻辑,访问行为留痕可查。
-- 通过快照元数据实现时间旅行,回查历史数据
SELECT * FROM iceberg_catalog.dws.ods_order
VERSION AS OF 8721345553
WHERE region = '华东';
-- 查看快照历史,定位数据变更点
SELECT snapshot_id, committed_at, operation, summary
FROM iceberg_catalog.dws.ods_order.snapshots
ORDER BY committed_at DESC LIMIT 10;落地过程中的典型问题与应对
第一个常见问题是小文件。流式作业checkpoint间隔通常在分钟级甚至秒级,每次提交都会产生新的数据文件,跑上几天后一张表可能堆出几十万个小文件,查询性能急剧下降。应对方案是在写入侧合理设置 checkpoint 间隔和文件大小阈值,配合定期的 compaction 任务做小文件合并,同时利用Iceberg的过期快照清理机制及时回收无效文件。
-- 定期触发压缩合并小文件
CALL iceberg_catalog.system.rewrite_data_files(
table => 'dws.ods_order',
strategy => 'binpack',
options => ('target-file-size-bytes' => '134217728')
);
-- 清理过期快照,保留最近7天
CALL iceberg_catalog.system.expire_snapshots(
table => 'dws.ods_order',
older_than => TIMESTAMP '2024-01-01 00:00:00',
retain_last => 100
);第二个问题是schema演进的兼容性管理。业务字段变更是常态,加字段一般安全,但改类型、删字段必须谨慎,尤其是流作业正在运行时变更schema,要确保新旧数据在同一个快照里能正确读取。建议建立schema变更的审批流程,任何DDL操作先在元数据中心登记,评估下游血缘影响后再执行。
第三个问题是资源隔离。批任务和流任务混部在同一个集群时,批作业的大规模shuffle可能挤占流作业资源,导致实时链路延迟抖动。生产环境建议通过Yarn队列或Kubernetes namespace做物理隔离,流作业独占一批资源保证稳定性,批任务用弹性资源按需伸缩。监控层面要同时关注流作业的checkpoint耗时、批作业的队列等待时长这两个关键指标,它们分别反映了两类任务的健康状态。
总体来看,批流一体的落地路径可以总结为三步:先用统一的计算引擎收敛开发模式,再用数据湖表格式统一存储与事务语义,最后把元数据中心建起来支撑长期治理。三步缺一不可,跳过元数据治理直接上线的团队,往往在半年后就会重新陷入口径混乱的老问题。