上个月帮一个做网约车综合数据分析的朋友做架构评审,他在要不要把订单指标从T+1改成实时推送这件事上纠结了很久。业务方天天催"我要看实时订单量",技术团队则担心实时链路搞下去会把数仓的底子搞乱。这其实不是我第一次碰到这种纠结了——几乎所有往数据平台方向走的团队,迟早都会卡在同一个问题上:到底哪些数据值得上实时ETL,哪些老老实实跑批处理就行。
我当时没有直接给结论,而是跟他把批处理ETL和实时ETL从头到尾掰开捋了一遍。今天这篇文章,就算是我这些年做数据架构选型的一个复盘吧。我没有打算把它写成教科书式的"实时好还是批处理好"的二选一,而是想把真正影响决策的那些东西——时延、成本、一致性、运维复杂度、团队现状——全部摊到桌面上,再用几个实际场景告诉你,我一般是怎么帮项目做这个判断的。
1. 批处理ETL的底牌:稳定性极强,但时延藏在调度窗口里
1.1 批处理ETL的典型工作流长什么样
先说说批处理ETL。这个名字听着好像很"传统",但你得承认,今天绝大多数数据仓库的骨架,依然是靠它撑起来的。
典型的工作流大概是这样的:业务数据库或者日志文件,在每天固定的时间点(通常是凌晨),通过Sqoop、DataX或者Spark任务,被拉取到数据仓库的ODS层。然后再通过一系列Hive SQL或者Spark SQL做清洗、过滤、脱敏、维度退化,一层一层从DWD跑到DWS,最后落到ADS层,供报表系统、BI工具查询。整个过程就像流水线上的夜间班次,凌晨开工,天亮之前把昨天的数据整理得整整齐齐。
我之前接触过一个网约车综合分析项目,用的就是这套思路:原始订单表和轨迹日志先落到Hive,然后用Spark做数据清洗——修正经纬度异常值、剔除重复订单、统一时间字段格式——清洗完的数据再喂给下一层的统计分析任务。整个链路清爽、可控,每个环节都有明确的输入输出,出了错可以按分区重跑,绝不会有"数据跑丢了但没人知道"的情况。
1.2 为什么说批处理的核心优势不是"快",而是"稳"
很多人以为批处理ETL的卖点是能处理海量数据,但我觉得它有两大真正的底牌:
第一是容错模型极其成熟。一个批处理任务挂在半路,只需要找到出错的依赖项,修完数据原地重跑,不会对下游产生任何影响。因为下游根本还没开始读这一批数据。这种"错了我可以重来"的底气,在实时链路里是奢侈品。
第二是成本可预估。批处理集群的负载是波动的——凌晨跑数时CPU飙高,白天可能闲置。你完全可以把它规划成夜间集中计算,资源利用率高,运维也好安排。在云上跑的话,弹性伸缩策略也简单粗暴:凌晨加节点,白天减节点。
它最大的问题也很直白:时延被调度窗口锁死了。经典T+1的数仓,你在今天下午3点问"今天上午10点到11点产生了多少新订单",答案大概率是"明天早上才能给你"。如果业务方需要的恰好是分钟级甚至秒级的数据,批处理这条路本身就走不通——这不是优化能解决的问题,是架构方向的限制。
这里我想补充一个很多人容易忽略的细节:批处理ETL的"定时"并不意味着它天生不能更频繁。实际上,如果调度系统支持,你完全可以每小时、每10分钟跑一个微批次任务。但这样做的代价是任务调度次数增多,近线处理的逻辑开始介入,数据质量校验、分区管理、副本一致性的维护成本都会明显上升。所以,批处理和实时的分界线,从来不是"任务跑得多频繁",而是"数据是否在产生后立刻被消费"。搞清楚这个概念以后,选型的很多纠结其实能少一半。
2. 实时ETL的本质:连续查询,而不是"跑得更快的批处理"
2.1 实时ETL的数据流模型:从Kafka到Flink的基本套路
聊实时ETL之前,我通常先给团队讲一句话:批处理是一次性跑完的查询,实时ETL是永远跑不完的查询。这两者在架构上的差异,远远大于"快和慢"。
批处理的数据是有限数据集,任务终有结束的一刻;而实时ETL处理的是无边数据流,数据源源不断进来,任务从启动那一刻起就一直在跑,直到你主动停止它。用这个视角去看技术选型,很多组件为什么要那样设计,就都说得通了。
最常见的实时ETL链路是:业务日志或数据库变更,通过Canal监听、Flume采集或者API回调,实时写入Kafka这类消息队列,接着由Flink这类流式计算引擎做实时清洗、关联、聚合,最后把结果写入Kafka、Elasticsearch、Doris、ClickHouse或者Redis,供下游的实时特征服务、数据大屏、业务告警去消费。
以我之前搭过的实时特征服务为例:用户每次点击、每次下单,事件会毫秒级进入Kafka;Flink消费后做属性拼接和窗口聚合——比如"过去5分钟这个品类被加购了多少次"——算出的特征写进KV存储,供上游的推荐或风控系统实时读取。整条链路的核心并不是"算得快",而是"事件一发生,特征就得更新"。
2.2 状态、水印与精确一次:实时ETL额外付的账
很多从批转实时的团队,第一个认知冲击都来自这几点。
状态管理。实时聚合天然要保存中间状态——比如"当前窗口内每个商品的累计销售额",这个状态不落盘就有丢失风险,落盘又会带来性能开销。Flink为了解决这个问题设计了状态后端(内存、RocksDB、文件系统等),还要求你为状态设置合适的TTL,不然状态无限膨胀,最终宕机。
水印与乱序。实时流里数据到达顺序经常和事件实际发生顺序不一致——网络延迟、客户端缓冲都可能造成乱序。假如你按"事件时间"做窗口统计,但不管乱序数据,统计结果就会偏。Flink用Watermark机制来处理这个问题,它是一个"时间截止线":告诉你系统等到某个时间点为止了,再晚到的数据要么丢弃、要么扔到侧输出流单独处理。这个机制非常实用,但我见过的几乎所有实时数仓团队,都在这上面交过学费——水印设置太保守,指标延迟严重;设置太激进,统计结果又经常被晚到数据"打脸"。
精确一次语义。消息从Kafka到Flink再到下游,任何一个环节都可能重复消费或漏消费。实时ETL要做到"精确一次"——即每条数据对结果的影响恰好一次——需要靠Checkpoint机制跨环节对齐提交位点。这套机制维护成本不低,一旦Checkpoint频繁失败,任务延迟和反压立刻就会上来。我见过不少团队,上线初期为了追求"标准配置",最后被Checkpoint超时搞得焦头烂额,后来老老实实退到"至少一次"语义,靠下游去重来兜底。
所以说,实时ETL的能力边界并不只在"快"这一个字上。后面拖着一大串状态管理、一致性策略、乱序处理的复杂度。这些就是你选择实时链路时必须提前预算好的运维成本和系统成本。
3. 选型决策框架:从业务时效到团队成本,五个维度做判断
3.1 时效需求:不是"越快越好",而是"在这个时延下数据是否有用"
每次聊选型,我会先问业务方一个很朴素的问题:你要的这个指标,超过多少分钟就没意义了?
这个问题的答案,基本上已经把大半条路给画出来了。我按自己的经验把时效需求分成三档:
| 时效档位 | 典型场景 | 建议链路 | 参考时延 |
|---|---|---|---|
| T+1档 | 经营报表、财务结算、监管报送 | 传统离线数仓(Hive/Spark) | 次日凌晨 |
| 分钟级档 | 运营看板、实时监控、销量统计 | 微批/近实时(Spark Streaming / Flink窗口任务) | 1~10分钟 |
| 秒级档 | 风控特征、实时推荐、设备告警、实时公交位置 | 实时流式链路(Kafka + Flink + KV/OLAP存储) | 秒级 |
你注意,中间这一档"分钟级"其实特别有意思:它经常被业务方喊着"我要实时",但刨根问底之后发现,5分钟延迟完全够用。这种时候你完全可以不做完整的实时链路,用Spark Structured Streaming按分钟跑微批更划算,资源和稳定性都好得多。
3.2 数据一致性、成本与团队能力:三个容易被低估的维度
除了时效,我至少还会看三件事。
数据一致性要求。如果你的指标最终要进财务报表,或者用于对账,那么实时链路的"近似准确"可能根本过不了审计这一关。批处理可以在跑完后对全量数据做精确计算,出了问题整分区重跑,这属于强一致方案。如果你的业务允许最终一致——比如大屏上的实时订单量,稍微偏差个几百单没人追究——那实时方案才可以进入备选池。
成本预算。实时链路的成本不只是多买几台机器这么简单。Kafka集群、Flink计算节点、状态存储、下游OLAP的副本、监控告警体系,每一项都是持续支出。而且流式任务7×24小时常驻,几乎没有"低峰期"可以缩容。我做预算的时候有个粗略的体感:同样的数据量,实时链路的计算+存储成本一般是批处理的2到3倍起步。如果你的数据量本身不大,但各部门都喊着要实时,先把数据湖或数仓的分层做扎实,可能比盲目上实时链路更划算。
团队能力。这个维度最容易被忽略,但踩坑概率最高。实时ETL的排查难度比批处理高一个量级:批处理任务失败,日志+重跑基本能解决问题;实时任务出问题,可能是Kafka堆积、Checkpoint超时、状态后端容量不足、乱序数据处理不平和……一大堆变量交织在一起。团队里如果没人能看懂Flink的反压监控面板,我建议先别急着全链路实时化。我见过太多项目,实时链路搭起来了,业务方也开始依赖了,结果一遇到数据倾斜就没人能修,最后还是要偷偷回批处理。先培养人,再上系统,比先上系统再找人救火,要稳妥得多。
3.3 整合成一个决策表:我平时怎么判断
把上面几个维度合到一起,我习惯用一个小决策表来落地:
| 判断条件 | 批处理 | 近实时/微批 | 全实时 |
|---|---|---|---|
| 数据时效要求 | T+1可接受 | 分钟级即可 | 秒级必需 |
| 数据一致性要求 | 强一致/可审计 | 最终一致即可 | 最终一致即可 |
| 成本敏感度 | 高 | 中 | 低 |
| 团队流式计算经验 | 无 | 有基本经验 | 较强排障能力 |
| 数据规模 | 大但平稳 | 中等 | 大且持续增长 |
这个表不是死的,但它能快速帮你排除明显不合理的选项。比如三个条件同时命中"批处理"那一列,就别再花时间调研实时方案了;反过来,如果业务明确要求秒级特征,那也别纠结"先跑个离线版看看效果"——直接小步快跑搭实时链路更实际。
4. Lambda还是Kappa:混合架构的真实代价
4.1 双链路合并:听起来简单,做起来容易翻车
聊到"既有实时又有批处理",那就绕不开Lambda架构。这套架构的思路很简单:同一份数据,用两条链路分别处理——批处理链路负责全量、精确的计算,实时链路负责低延迟的近似结果,最终由服务层把两条链路的结果合并返回给业务。
听起来很合理对不对?但我实际接触过好多个用Lambda架构维护了很久的项目,几乎无一例外都会在"合并"这一步上持续出血。
最典型的问题是数据对不上。同一个指标,实时链路在凌晨12点整点算出的当天结果,和批处理链路早上8点重新算出来的结果,大概率会有出入。原因很多:实时链路可能用了窗口近似聚合,批处理是全量精确计算;乱序数据在实时链路被丢弃了一部分,批处理却能完整看到。于是在业务侧就出现了经典的"昨晚实时看板显示100万,今天早上报表出来却是98万"的尴尬。
这个差异一旦出现,排查成本非常高,因为要同时查两条链路的处理逻辑、数据边界、计算语义。更麻烦的是,两条链路的代码往往是两套完全不同的技术栈——一套Spark SQL、一套Flink SQL——逻辑哪怕再努力对齐,细节上也很难完全一致。维护了三个月以后,团队基本就处于"每天在给实时链路打补丁,让它的结果尽量逼近批处理"的状态。
所以我对Lambda架构的立场是这样的:如果你家有强大的数据服务层,且两个结果之间的差异能够被业务容忍或解释清楚,那Lambda是没问题的;但如果合并层的代码只是简单粗暴地把两路结果加在一起,这家架构迟早会变成你团队的一个慢性病。
4.2 Kappa架构与湖仓一体:今天还有更轻的解法吗
Kappa架构主张只用一条实时链路,通过Kafka一类的日志系统来保存历史数据。需要重新计算历史时,不重写批处理代码,而是把Kafka的消费位点拨回去,用同一套实时代码"重新跑一遍"。这个想法很优雅,但早期落地有个硬伤:Kafka保留全量历史数据的成本太高,重放太慢。
不过这几年湖仓一体技术发展以后,情况变了不少。像Hudi、Iceberg、Paimon这类支持行级更新的数据湖格式,给了我们一个新的中间选项。你可以在实时链路上用Flink直接写Hudi/Iceberg表,既能按微批的方式落文件,又能支持行级Upsert;下游查数的时候,得到的是一个"既接近实时、又支持回溯修正"的结果。这在很大程度上把Lambda和Kappa的优点做了折中。
我最近一个做实时公交API的项目就是这么落地的:车辆位置事件进Kafka,Flink做实时清洗和路径匹配,结果同时写Kafka(供毫秒级查询)和Hudi表(供历史回放与分析)。历史指标要重算时,不再需要一条独立的批处理链路,只要重放Hudi里的明细数据就行。坦白讲,这种方案在小团队、中等数据体量项目里,落地比双链路Lambda舒服得多。
5. 三类典型场景的落地拆解:从网约车到实时大屏
5.1 网约车数据分析:哪些环节必须实时,哪些环节别碰实时
网约车应该是"实时+离线混合"最好的样本了。一个完整的网约车数据体系里,数据可以分成好几类,选型策略完全不同。
订单主状态链路(下单、接单、司机到达、行程开始、行程结束、支付)——这是典型的必须实时的场景。因为产品端要在用户点开APP的瞬间展示附近车辆实时状态,运营端要实时看到各区域订单密度来调度运力。这类数据走Kafka+Flink实时清洗,写实时宽表或者直接推送到查询服务,几乎是刚需。
轨迹轨迹和里程数据——这是最有意思的一类。轨迹数据量巨大,但上游业务对它的实时要求并不是全局性的:司机端行程进行中需要实时位置来做导航和计价预估值,但离线分析(热力图、常用路线挖掘、里程统计)完全可以放到批处理里去算。我在之前的网约车项目里看到过一本经典做法:轨迹落Kafka走实时链路给在线地图引擎,同时落一份Hive表给凌晨的Spark清洗和分析任务用。
财务和计价数据——从我的经验看,计价结算这类涉及钱的数据,我强烈建议不要为了"实时看数"把它做成全实时的核心链路。计价本身可以有实时预估(给司机端显示一个大概范围),但真正的账单生成和对账,老老实实走批处理,让它可以重跑、可审计。这是我反复跟业务方强调的一条底线:"实时"和"钱"接触越深,出风险的代价越大。
5.2 实时大屏与实时公交API:秒级更新的"伪实时"真相
很多人一想到数据大屏,就觉得背后一定是一套极复杂的实时计算平台。实际上我拆过不少大屏项目——包括用Flask+ECharts做的那种经典组合——它们的"实时"程度往往没你想象得那么高。
大屏上展示的订单量、用户量、销售额,如果刷新粒度为秒级,其实大部分数据并不需要从OLAP里做实时聚合查询。一个更稳的套路是:
- 用Flink做实时聚合,把每分钟的结果写进Redis或者Doris的预聚合表。
- 大屏后端(比如Flask接口)只负责从这些预聚合结果中做毫秒级查询。
- 前端ECharts每2秒轮询一次后端接口,刷新图表数据。
这套设计的本质是把"实时计算"和"实时展示"解耦。真正承担实时计算压力的只有Flink那一条链路,而后端接口和大屏前端的负担都很轻,也不容易被查询打垮。
实时公交API也是同样的套路。公交车的GPS位置上报到平台,实时链路负责计算"某辆车当前离某站还有多远、预计多久到站",结果写进Redis。用户通过APP查询时,其实是在读Redis里那份已经算好的结果,而不是让查询引擎现场实时测算。一边是后台永不停歇的计算,一边是前端毫秒级的读取——真正的实时应用,背后几乎都是这个"计算与展示分离"的模式。
6. 实测中的翻车记录:实时ETL最常见的三个坑
6.1 反压失控:Kafka堆积与Checkpoint超时是怎么拖垮整条链路的
如果你要上实时ETL,我建议先把下面这份"翻车清单"贴在显示器旁边。
第一个高频事故是Kafka消费堆积。表现为消费者组Lag持续走高,业务指标越来越旧,但任务本身没报错。最常见的原因是Flink里某个算子处理能力跟不上——比如一个需要远程调用的维度补全函数,单条数据处理耗时从1毫秒变成了50毫秒,吞吐直接垮掉。Flink的反压监控面板里能看到算子积压,但很多人上线后根本不看这个指标。
第二个高频事故是Checkpoint超时导致任务重启。状态过大、依赖的外部系统响应太慢、或者网络抖动,都会让Checkpoint迟迟不完成,最终触发Failover。重启后的任务要重新对齐状态,吞吐又掉下去,形成恶心循环。踩过几次坑以后我现在有个习惯:任何实时任务上线前,先做一次带状态大小的故障演练,确认重启恢复时间在可接受范围内,再放量接真实数据。
第三个高频事故是下游写入瓶颈。实时任务计算结果要写Elasticsearch或Doris时,下游索引刚好在做合并或建表,写入性能剧烈抖动,实时链路跟着被拖死。这事儿的本质是实时任务对下游的冲击是持续的、高频的,和批处理"写完就走"完全不一样。所以现在我做实时方案,一定会提前给下游设好写入QPS上限和缓冲策略,把"突发写入"转成"匀速写入"。
6.2 排障思路与止损心法:从"责任链"看问题
实时链路排障,我有几个固定套路。
一看消费延迟:如果Kafka的Lag在涨,优先在Flink任务侧找原因,不要先去查下游业务。二看算子反压:Flink面板上哪个算子压力大,就从哪个算子的逻辑入手,要么优化函数、要么调并行度、要么加资源。三看状态大小:状态异常膨胀,先查有没有忘记给状态设置TTL,再看是不是聚合维度爆炸了。
还有一个止损心法我想多聊两句:永远给实时任务留一个"手工重放"的开关。也就是说,Kafka里的原始数据按照时间戳保留,并允许你随时把某个Flink任务的消费位点拨回去重新消费。这样当计算结果出现偏差时,不用重新跑批处理任务,只要把任务停掉、位点回拨、再启动,就能重新计算指定时间段内的数据。Kafka的默认保留时长一般是7天,如果你有对账或者重算需求,建议把保留时长调到15到30天——这点存储成本,和"数据错了无法重算"的代价相比,实在便宜太多了。
我在实际项目里还有一个原则:实时链路刚上线时,一周内不砍掉旧的批处理任务。让两条链路并行跑,每天对账。只有当实时结果的准确率连续一周达标,才允许逐步下线批处理任务。这不是保守,而是给新系统留一个退路——你永远不知道实时链路会在第几天暴露出什么奇奇怪怪的边界问题。
7. 写在最后的建议:让选型逻辑回归"业务问题"
扯了这么多,最后还是想回到选题本身。
我在做技术咨询的时候,最怕听到的一句话是"我们数据量很大,所以我们要上实时"。数据量大小,从来不是实时ETL的充分条件。真正决定该不该上实时ETL的,是业务是否真的产生了对低延迟数据的需求:用户要在大屏看到秒级变化的趋势,风控系统要在交易发生时立刻拿到特征,运营要在流量异常的瞬间收到告警——这些需求才配得上实时链路的成本和复杂度。
如果仔细盘点一下,你会发现大多数业务场景对数据时延的要求其实落在"分钟级"档位,而这一档用微批或者近实时方案往往就能覆盖得很好,完全没必要一上来就上全流式架构。反过来,真正需要秒级实时特征的场景,通常就那么一两个核心链路,那就把那部分做好、做精,而不是让整个数仓都跟着实时化。
我自己的习惯是:先用最便宜的方案把业务跑通,再根据真实的时效需求决定要不要加实时链路——大多数情况下,最便宜的那条路已经能解决80%的问题了。剩下的20%,等它变成了业务的一等公民,再认真投入不迟。数据架构没有绝对的对错,只有"在这个时间点、这个资源约束下最合适的选择"。