☰
ETL选型实战:实时与批处理怎么选?数据工程师避坑指南
2026/10/2 8:44:59 网站建设 项目流程

做数据平台这些年,我接手过不少中大型项目,发现一个特别有意思的现象:项目启动会上大家讨论得最激烈的,往往不是数据模型怎么设计,也不是指标口径怎么定,而是ETL到底用实时还是批处理。这个选择题如果没想透,后面建仓库、写管道、调任务的时候,到处都是返工的坑。实时ETL和批处理ETL从来不是“谁替代谁”的关系,它们面对的数据形态、业务诉求、成本结构完全不一样。这篇文章我会把两种模式的底层逻辑、适用边界、选型策略,以及我在真实项目里踩过的坑一次性讲清楚,适合正在做离线数仓、实时数仓的工程师,也适合刚入行想建立整体认知的数据同学参考。

1. 先把ETL这件事想清楚:它到底在解决什么问题

1.1 一张图看懂ETL在大数据链路里的位置

ETL是Extract(抽取)、Transform(转换)、Load(加载)的缩写,三个动作合起来,就是把散落在各个业务系统里的原始数据,经过清洗、加工、标准化后,搬到数据仓库、数据湖或者分析型数据库里供下游使用。你可以把它想象成一个城市自来水系统:业务数据库是水源地,数仓是蓄水池,ETL就是中间的水处理厂——原水不能直接入户,必须经过沉淀、过滤、消毒才能进到用户的水龙头。ETL干的就是这件事,只不过它处理的是数据和各种业务逻辑。

在大数据链路里,ETL是整个体系的起点和生命线。上游是业务系统产生的日志、订单、用户行为、设备上报等数据,下游是BI报表、用户画像、推荐系统、实时大屏、风控模型。ETL做得好不好,直接决定下游分析结果的准确性和及时性。我见过很多项目,建模画了一堆漂亮的图,结果ETL管道三天两头跑挂,数据对不上,最后整个数仓沦为摆设。

1.2 批处理和实时的本质差异:不是快慢问题,是处理范式问题

很多人觉得批处理和实时的差异就是“快和慢”,这其实是最大的误解。它们的本质区别在于处理的数据形态不同。

批处理ETL处理的是“已经发生完”的数据。比如凌晨2点开始跑昨天一整天的订单数据,这时候当天所有订单都已经入库了,数据是完整的、有边界的。它面对的是一个有限的数据集,所以可以反复重算、全量扫描、做复杂的多表关联,也不用担心数据突然中断或者迟到。

实时ETL处理的是“正在发生”的数据。订单下了、位置上报了、日志产生了,数据在几毫秒内就到了Kafka,这时候你要决定马上处理还是攒一批再处理。数据是无限的、无序的、可能延迟的,窗口边界要自己定义,数据丢了要能找回来,消息重复了要去重。这套复杂度是批处理完全没有的。

所以我经常跟团队说:选型的第一步不是看技术栈,是看你手里的数据在时间维度上是什么形态。你是在跟一段完整的历史打交道,还是在追一条不断向前流动的河。

1.3 从数仓分层看两种模式的落脚点

大数据架构通常分成四层:ODS(原始数据层)、DWD(明细数据层)、DWS(汇总数据层)、ADS(应用数据层)。批处理和实时ETL在这四层里的角色完全不同。

在ODS层,批处理和实时各自维护一套原始数据副本——离线副本落在Hive表里,实时副本落在Kafka的topic里。到了DWD层,批处理做清洗去重、维度补充、统一字段格式,产出标准化的明细表;实时的DWD层则通过流式任务做同样的动作,只是用窗口和状态来代替批处理的全局扫描。DWS层是两者的分歧点:批处理用Spark SQL做明细表的聚合,产出T+1的汇总表;实时链路则用Flink做持续聚合,把结果不断更新到Redis、Doris或ClickHouse里,支撑秒级的大屏和看板。

理解了这个分层逻辑,你才能明白:不是所有环节都需要实时,也不是所有环节都能接受批量延迟。有些层可以共享一套管道,有些层必须拆开跑,这是后面所有选型讨论的基础。

2. 批处理ETL:为什么它仍然是数据平台的中流砥柱

2.1 技术栈盘点与选型:Hive、Spark SQL、调度框架怎么搭

现在主流的批处理ETL技术栈,基本是以Hive和Spark SQL为核心。Hive适合超大吞吐量的离线清洗,把SQL翻译成MapReduce或Tez任务跑在YARN上;Spark SQL则凭借内存计算的优势,在同样的数据量下往往比Hive快数倍,而且能无缝衔接Scala、Python写的自定义处理逻辑。

调度层常见的是Apache DolphinScheduler和Airflow。两者各有侧重:DolphinScheduler对大数据生态更友好,原生支持Hive、Spark任务类型,带可视化的工作流编排和补数机制;Airflow胜在Python生态强,适合和机器学习管道打通。我自己在纯数仓项目里用DolphinScheduler更多,因为它的补数、重跑、依赖触发做得很直接,能减少很多运维沟通成本。

还有一点要注意:批处理的数据存储选型,建议用支持列式和分区裁剪的方案,比如Hive的ORC/Parquet格式加上合理的分区策略。分区字段选错了,后面全表扫描的代价会让你哭都哭不出来。最常见的做法是按天分区,数据量大的场景再叠加按小时分区或者按业务ID哈希分桶。

2.2 一个典型的批处理ETL任务长什么样

我拿一个网约车订单日表来举例,这是个非常经典批处理场景。每天凌晨,调度任务开始执行:第一步从ODS层读取前一天的全量订单日志,第二步过滤掉状态异常的脏数据(比如订单金额为负、起终点经纬度都在海外的乱码数据),第三步把经纬度映射到城市和商圈,同时关联司机维表和乘客维表补上姓名、车型等信息,第四步按城市、时段、订单状态做多维聚合,最后覆盖写入DWS层的日汇总表。

对应到代码上,核心逻辑一般长这样:

INSERT OVERWRITE TABLE dws_order_day_agg PARTITION(dt='${bizdate}') SELECT city_id, hour_bucket, order_status, COUNT(*) AS order_cnt, SUM(order_amount) AS gmv, AVG(order_amount) AS avg_amount FROM dwd_order_detail WHERE dt='${bizdate}' AND order_amount >= 0 AND city_id IS NOT NULL GROUP BY city_id, hour_bucket, order_status;

这一条SQL看着简单,实际跑的时候会暴露出各种问题。比如某一座城市订单量特别大,Group By的时候就会产生数据倾斜,所有数据压到同一个Reduce上,其他节点跑完干等。这种问题我在第五节会专门展开讲。

2.3 批处理的三大困境和工程化解法

第一个困境是延迟天花板。T+1的节奏,意味着今天的数据明天早上才能看到。业务方如果要做当日实时监控或快速决策,批处理根本接不住。解法不是把批处理调快,而是把真正需要实时的场景剥离出去,让批处理和实时各管一段。第二个困境是计算成本高。批处理跑大表Join和全量重算,动辄几十亿行数据,每跑一次都是资源的真金白银。常见的化解法是增量抽取加增量计算,而不是每天全量重跑;同时用数据血缘分析,避免下游任务串联式的重复计算。第三个困境是数据倾斜,这个我在实操部分会详细说,核心套路就是加盐拆分、两阶段聚合,或者用Skewed Join优化。

批处理现在依然是数据平台的中流砥柱,原因很简单:它稳定、可重试、结果可解释。对于报表、财务对账、月度分析这类场景,你不需要秒级数据,你需要的是“无论几点跑,结果都是准的”这个保障。批处理在这方面天然可靠。

3. 实时ETL:低延迟的红利与技术代价

3.1 微批处理和真流式到底差在哪

实时ETL领域有一个常年争论:Spark Streaming的微批(Micro-Batch)和Flink的真流式(True Streaming)到底怎么选。理解这个问题的关键,在于“批”这个字。

Spark Streaming把实时数据切成一小段一小段的小批,每隔几秒或几十秒提交一次,本质上还是批处理,只是批变小了。它的好处是架构简单、和Spark生态天然融合,出问题时可以精确地重试某一个小批。坏处是延迟存在下限,延迟抖动比较明显,而且窗口计算和状态管理的能力偏弱。

Flink做的是真正的逐条流式处理,事件到了就处理,延迟可以压到几百毫秒甚至更低。它在窗口计算上非常灵活——滚动窗口、滑动窗口、会话窗口都能直接表达,还内置了精确一次(Exactly-Once)语义的状态管理。代价是学习曲线陡峭、调试复杂、状态后端管理麻烦,资源占用也比微批高。

我的个人建议是:如果只是做实时大屏或者简单过滤转发,Spark Streaming够用;如果要做复杂的事件驱动、乱序数据处理、双流Join,直接上Flink,省得后期重写。

3.2 Lambda架构和Kappa架构:实时数仓的两条路线

做实时ETL,最难的不是写Flink任务,而是整体架构怎么搭。这里绕不开两个经典架构模式。

Lambda架构是“批处理层+实时速度层”双轨制:同一份数据,离线链路跑出准确的批量结果,实时链路跑出低延迟的近似结果,服务层再按需合并。这种架构的好处是两条链路互不干扰,每条都用自己擅长的方式处理数据;坏处是要维护两套代码、两套管道,口径稍微不一致就会出现离线数和实时数对不上的问题。

Kappa架构则干脆只用一条实时流,把数据全部通过Kafka等消息队列流式处理,需要重算时重新回放历史数据。它省掉了一套维护成本,但对实时系统的稳定性、状态大小、数据回溯能力要求极高。我见过不少团队选Kappa,结果一次历史数据回放就压垮了Kafka集群。

实际项目里,真正落地最多的是改进版的Lambda——离线链路继续用Hive/Spark跑T+1的精确结果,实时链路只负责产出T+0的快速指标,两边在服务层各取所长。这是我最推荐的务实路线,没有之一。

3.3 实时ETL特有的三个难题:一致性、状态管理、数据质量

实时ETL真正难的是下面三个问题,个个都能让你深夜被电话叫醒。

数据一致性方面,消息队列的At-Least-Once投递会导致重复数据,Flink虽然提供Checkpoint机制实现精确一次,但是Sink端如果写的是Redis或MySQL这类不支持事务的系统,依然可能出现重复写入。解决办法通常是下游做幂等,比如用唯一键去重,或者在写入时带上业务时间戳做覆盖写。

状态管理方面,Flink做聚合、去重、双流Join都依赖状态。状态太大、Checkpoint太频繁,都会拖垮任务;状态后端选不好,还会频繁Full GC。通用的做法是大状态场景用RocksDB状态后端,配合增量Checkpoint,同时合理设置TTL清理过期状态。

数据质量方面,实时链路处理的是无边界数据,乱序、迟到、缺失是家常便饭。Watermark设置太保守,结果迟迟不触发计算;设置太激进,又会有大量迟到数据被丢掉。我的建议是:核心指标宁可稍微等一等,等Watermark延迟3到5秒,也不要用早到的Watermark造成数据缺失。

4. 选择策略:用六个维度搞定绝大多数项目的选型

4.1 六个评估维度详解

每次有团队跑来问我“到底选实时还是批处理”,我基本会用下面这六个维度帮他们过一遍,最后答案自己就浮出来了。

**维度一:延迟SLA。**这是最硬的约束。业务方要的是“明天看到昨天的报表”,还是“3秒内看到当前的实时订单量”?前者批处理完全够,后者必须引入实时链路。延迟要求每降低一个量级,技术复杂度都是成倍增加的,所以一定要跟业务方确认清楚。很多所谓的“实时需求”仔细一问,其实30分钟内的延迟都能接受,那完全可以考虑用微批或者准实时方案。

**维度二:数据量级与峰值。**每天几千万条和每天几十亿条的处理策略完全不同。批处理的成本瓶颈在计算资源,实时的成本瓶颈在存储和状态管理。数据量越大,实时链路对状态后端和Kafka的带宽要求越苛刻,成本指数级上升。建议先做数据量摸底,再决定是否值得上实时方案。

**维度三:成本预算。**实时集群的规格、Kafka的存储副本、Flink的并行度、下游OLAP的机器,每一项都比离线方案烧钱。我见过一个项目,实时链路成本是原来离线链路的三倍,最后算下来业务收益完全覆盖不了。成本评估一定要在技术选型之前做,而不是做完再补。

**维度四:团队技能栈。**如果团队只熟悉Hive和SQL,强行上Flink意味着数周的培训和数个月的踩坑期。反之,如果团队对流计算有经验,批量那套反而觉得不够灵活。选型一定要考虑团队现状,而不是招聘JD上写什么就选什么。

**维度五:数据一致性要求。**如果下游是财务对账,任何一个数字错了都要出大事,这种场景优先选批处理,因为它可以反复重算、交叉验证。如果下游是实时风控、实时推荐,允许一定的近似误差,实时链路才敢放开用。

**维度六:下游消费形态。**报表、月报、T+1分析适合批处理;实时大屏、实时告警、实时特征拼接,必须走实时。还有一个常见的中间态:下游想要“分钟级延迟的榜单”,这个用Spark Streaming的微批就能搞定,不一定要上Flink。

4.2 三种典型架构形态怎么定

把六个维度过完,大部分项目的架构形态也就清楚了。

第一种是纯批处理架构,适合报表型数仓、财务分析、传统BI。所有ETL任务按天或按小时调度,数据全部落到Hive/Iceberg类存储,下游查询走Spark SQL或Presto。优点是稳定、成本可控、维护简单;缺点是不能支撑实时场景。

第二种是纯实时架构,适合风控、实时反欺诈、IoT监控等对延迟极其敏感的场景。所有数据直接进Kafka,Flink任务持续处理,结果实时写入OLAP或在线存储。优点是一条链路响应极快;缺点是搭建门槛高、成本高、回溯困难。

第三种是混合架构(改良Lambda),适合大多数综合性业务平台,既出日报也要实时大屏。离线链路负责准确计算和补数据,实时链路负责秒级指标和监控,两条链路在DWS层对接,通过口径定义保证两边数据能对上。这是我最推荐的方案,后面网约车案例就是按这个思路做的。

4.3 我常用的决策流程与话术

我一般用三步走的方式帮团队做选型决策。第一步,把所有业务的延迟需求列成一张表,按“秒级-分钟级-小时级-T+1”分组,这一步能筛掉一半不靠谱的实时需求。第二步,对留下来的秒级和分钟级场景做数据量评估,估算Kafka吞吐、Flink并行度和存储增量,得出一个粗略的成本数字。第三步,把成本和延迟收益摆到业务方面前,用一句话收尾:“能用批解决的别上实时,实时只留给那种真的等不了的业务。”

这套流程听起来简单,但很管用。因为它逼着业务方把自己的需求具体化,而不是泛泛地说“我们要实时化”。一旦进入具体的SLA和成本讨论,很多伪需求自然就消失了。

5. 实操复盘:一个网约车数据项目的ETL落地全过程

5.1 场景设定:离线加实时双链路怎么切分

我拿之前做过的一个网约车数据平台举例,这个项目比较有代表性,既涉及高吞吐的订单流水,又有对延迟敏感的司机位置轨迹。我们的切分逻辑是:订单事实表走离线Hive链路,做T+1的城市、时段、司机维度的详细分析;司机实时位置和订单状态变化走Kafka+Flink链路,驱动实时大屏和司乘调度预警。

离线链路的核心是Hive数仓分层。ODS层直接落Kafka同步过来的订单日志,DWD层做数据清洗——过滤金额异常的订单、解析经纬度、处理时区问题、关联维表。DWS层按城市和小时聚合。ADS层供BI查询。调度用DolphinScheduler,每天凌晨两点自动跑,跑完自动发数据质量报告。

实时链路的核心是Kafka承接订单状态变更事件,Flink消费后做过滤清洗、维表关联、按城市维度滚动聚合,结果写到Redis,Flask后端接上ECharts做实时大屏展示。为了不让实时链路过度占用资源,我们只保留了一个大屏场景和两个预警规则,剩下的实时需求全部砍到了分钟级准实时方案里。

5.2 数据质量检查框架的搭建要点

这个项目让我印象最深的不是Flink任务写得多流畅,而是花了大功夫搞数据质量检查框架。刚开始的时候,离线跑完和实时计算结果经常差几个百分点,排查起来极其痛苦。后来我们把检查做成了自动化规则引擎,核心覆盖五个维度。

完整性检查,重点看主键是否为空、必填字段的非空率是否达标。准确性检查,比如订单金额不能为负数、经纬度必须在合理范围内、时间戳不能是未来时间。一致性检查,同一订单在明细表和汇总表里的金额必须一致,不同表里同一个城市ID对应的城市名必须相同。及时性检查,重点监控数据从产生到落库的耗时,判断管道是否存在堆积。唯一性检查,通过统计订单ID的重复率,判断是否有重复写入。

每个检查规则都配置了阈值和告警级别,一旦超过阈值自动通知对应负责人。这套框架上线后,数据问题平均发现时间从以天计缩短到了小时级,甚至在问题还没被业务感知到的时候就已经被我们拦住了。

5.3 数据倾斜等高频问题的现场处理记录

项目里有几个让我印象深刻的坑。第一个是离线大表Join时出现的数据倾斜。当时跑一个订单明细和司机维度关联的任务,每天凌晨都在某个时段卡死,后来定位到原因:一座超大城市贡献了将近30%的订单量,这一批数据全部落到同一个Reduce上。解决方法是把大Key加上随机盐前缀,先做一轮局部聚合,再去掉盐做全局聚合,也就是两阶段聚合,任务从40分钟直接降到了12分钟。

第二个是Hive小文件过多导致NameNode压力大。由于下游有大量的按天增量写入,产生的文件数量爆炸,我们手动加了合并参数:在Spark作业中开启了Hive的合并输出功能,把小文件合并成128MB级别的文件,NameNode的堆内存压力立刻降了下来。

第三个是实时链路的反压问题。Flink任务偶尔会出现Kafka消费延迟飙升,排查后发现问题出在下游Redis写入出现性能瓶颈,消费速度跟不上生产速度。优化方案是加大Sink端的并行度,同时把大批量写入改成分批异步写入,消费延迟很快就恢复了正常。

6. 常见问题速查表与避坑清单

6.1 高频问题排查速查表

我整理了批处理和实时ETL项目里最常遇到的10个问题,以及对应的排查思路,做成一张速查表,方便你直接照着排查。

问题现象常见原因排查与解法
批处理任务越跑越慢数据倾斜或小文件过多检查Spark UI的Stage耗时分布,用加盐两阶段聚合,开启小文件合并
Hive查询查询半天不出结果未走分区裁剪,全表扫描检查执行计划是否命中分区,补充分区字段过滤条件
实时任务消费延迟持续上涨下游Sink瓶颈或状态过大检查反压指标和Kafka消费Lag,增加并行度或调整状态后端
离线数据和实时数据对不上口径不一致或窗口边界不同统一指标口径,确认Watermark和窗口边界,增加对账任务
数据重复导致聚合翻倍消费端未做幂等下游使用唯一键去重,开启Flink的精确一次语义
F link任务频繁重启Checkpoint超时或OOM调整Checkpoint间隔,改用RocksDB状态后端,排查大状态算子
经纬度字段大量异常上游SDK上报格式问题在DWD层做规则过滤,不合法数据进异常表用于反馈上游
订单主键重复率高多数据源串联写入无去重建立主键去重规则,用状态或全局表做唯一性约束
维表更新后历史数据不一致维度缓慢变化未处理引入拉链表或者应用SCD策略
任务补数时依赖关系乱了调度依赖配置疏漏用DolphinScheduler的补数功能,理清上下游依赖再跑

6.2 三条独家避坑经验

第一条,别把实时链路当成万金油。很多业务的后台看板根本就不需要秒级刷新,用批处理加上一个30分钟的调度就足够了,没必要为了一个“看起来实时”的效果付出加倍的成本。第二条,成本评估要算全量成本。很多人只看实时集群的机器成本,忽略了Kafka的存储成本、Flink的并发资源、下游OLAP集群的扩容成本,这些加在一起往往超出预算。第三条,切换之前必须做一段时间的双跑验证。新链路上线后,先和旧链路并行跑两周,逐日对账,确认两边数据基本一致后再切换,这是避免数据事故最有效的手段。

7. 最后想说的话

回顾这些年做过的数据平台项目,我认为ETL选型没有标准答案,只有基于自身情况的权衡。批处理的可靠性和成本优势不会消失,实时的响应能力也确实是越来越多业务的刚需,两者在可预见的未来都会长期共存。

按照我的经验,团队成长和数据平台演进可以走一条比较稳的路径:先靠批处理链路把数仓、指标口径、数据质量体系打扎实,这是地基,一定不能省;等到业务量上来、团队对离线链路有了足够掌控力,再挑一两个真正等不了的场景切入实时,逐步扩大实时链路覆盖范围。这样做的好处是,每一次技术演进都是被业务需求驱动的,而不是为了技术而技术,团队踩的坑少,成本可控,业务方的满意度也更高。

最后再分享一个小技巧:数据平台建设初期就要把数据血缘管理做起来。不管是跑批任务还是实时任务,每一个指标、每一张表的上游依赖都清清楚楚的,当数据出问题的时候你才知道该去查哪一条链路。这个习惯越早建立,后面救你命的次数就越多。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询