尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

批处理与实时ETL选型指南:从时延、成本到架构决策

批处理与实时ETL选型指南:从时延、成本到架构决策 上个月帮一个做网约车综合数据分析的朋友做架构评审他在要不要把订单指标从T1改成实时推送这件事上纠结了很久。业务方天天催我要看实时订单量技术团队则担心实时链路搞下去会把数仓的底子搞乱。这其实不是我第一次碰到这种纠结了——几乎所有往数据平台方向走的团队迟早都会卡在同一个问题上到底哪些数据值得上实时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飙高白天可能闲置。你完全可以把它规划成夜间集中计算资源利用率高运维也好安排。在云上跑的话弹性伸缩策略也简单粗暴凌晨加节点白天减节点。它最大的问题也很直白时延被调度窗口锁死了。经典T1的数仓你在今天下午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供下游的实时特征服务、数据大屏、业务告警去消费。以我之前搭过的实时特征服务为例用户每次点击、每次下单事件会毫秒级进入KafkaFlink消费后做属性拼接和窗口聚合——比如过去5分钟这个品类被加购了多少次——算出的特征写进KV存储供上游的推荐或风控系统实时读取。整条链路的核心并不是算得快而是事件一发生特征就得更新。2.2 状态、水印与精确一次实时ETL额外付的账很多从批转实时的团队第一个认知冲击都来自这几点。状态管理。实时聚合天然要保存中间状态——比如当前窗口内每个商品的累计销售额这个状态不落盘就有丢失风险落盘又会带来性能开销。Flink为了解决这个问题设计了状态后端内存、RocksDB、文件系统等还要求你为状态设置合适的TTL不然状态无限膨胀最终宕机。水印与乱序。实时流里数据到达顺序经常和事件实际发生顺序不一致——网络延迟、客户端缓冲都可能造成乱序。假如你按事件时间做窗口统计但不管乱序数据统计结果就会偏。Flink用Watermark机制来处理这个问题它是一个时间截止线告诉你系统等到某个时间点为止了再晚到的数据要么丢弃、要么扔到侧输出流单独处理。这个机制非常实用但我见过的几乎所有实时数仓团队都在这上面交过学费——水印设置太保守指标延迟严重设置太激进统计结果又经常被晚到数据打脸。精确一次语义。消息从Kafka到Flink再到下游任何一个环节都可能重复消费或漏消费。实时ETL要做到精确一次——即每条数据对结果的影响恰好一次——需要靠Checkpoint机制跨环节对齐提交位点。这套机制维护成本不低一旦Checkpoint频繁失败任务延迟和反压立刻就会上来。我见过不少团队上线初期为了追求标准配置最后被Checkpoint超时搞得焦头烂额后来老老实实退到至少一次语义靠下游去重来兜底。所以说实时ETL的能力边界并不只在快这一个字上。后面拖着一大串状态管理、一致性策略、乱序处理的复杂度。这些就是你选择实时链路时必须提前预算好的运维成本和系统成本。3. 选型决策框架从业务时效到团队成本五个维度做判断3.1 时效需求不是越快越好而是在这个时延下数据是否有用每次聊选型我会先问业务方一个很朴素的问题你要的这个指标超过多少分钟就没意义了这个问题的答案基本上已经把大半条路给画出来了。我按自己的经验把时效需求分成三档时效档位典型场景建议链路参考时延T1档经营报表、财务结算、监管报送传统离线数仓Hive/Spark次日凌晨分钟级档运营看板、实时监控、销量统计微批/近实时Spark Streaming / Flink窗口任务110分钟秒级档风控特征、实时推荐、设备告警、实时公交位置实时流式链路Kafka Flink KV/OLAP存储秒级你注意中间这一档分钟级其实特别有意思它经常被业务方喊着我要实时但刨根问底之后发现5分钟延迟完全够用。这种时候你完全可以不做完整的实时链路用Spark Structured Streaming按分钟跑微批更划算资源和稳定性都好得多。3.2 数据一致性、成本与团队能力三个容易被低估的维度除了时效我至少还会看三件事。数据一致性要求。如果你的指标最终要进财务报表或者用于对账那么实时链路的近似准确可能根本过不了审计这一关。批处理可以在跑完后对全量数据做精确计算出了问题整分区重跑这属于强一致方案。如果你的业务允许最终一致——比如大屏上的实时订单量稍微偏差个几百单没人追究——那实时方案才可以进入备选池。成本预算。实时链路的成本不只是多买几台机器这么简单。Kafka集群、Flink计算节点、状态存储、下游OLAP的副本、监控告警体系每一项都是持续支出。而且流式任务7×24小时常驻几乎没有低峰期可以缩容。我做预算的时候有个粗略的体感同样的数据量实时链路的计算存储成本一般是批处理的2到3倍起步。如果你的数据量本身不大但各部门都喊着要实时先把数据湖或数仓的分层做扎实可能比盲目上实时链路更划算。团队能力。这个维度最容易被忽略但踩坑概率最高。实时ETL的排查难度比批处理高一个量级批处理任务失败日志重跑基本能解决问题实时任务出问题可能是Kafka堆积、Checkpoint超时、状态后端容量不足、乱序数据处理不平和……一大堆变量交织在一起。团队里如果没人能看懂Flink的反压监控面板我建议先别急着全链路实时化。我见过太多项目实时链路搭起来了业务方也开始依赖了结果一遇到数据倾斜就没人能修最后还是要偷偷回批处理。先培养人再上系统比先上系统再找人救火要稳妥得多。3.3 整合成一个决策表我平时怎么判断把上面几个维度合到一起我习惯用一个小决策表来落地判断条件批处理近实时/微批全实时数据时效要求T1可接受分钟级即可秒级必需数据一致性要求强一致/可审计最终一致即可最终一致即可成本敏感度高中低团队流式计算经验无有基本经验较强排障能力数据规模大但平稳中等大且持续增长这个表不是死的但它能快速帮你排除明显不合理的选项。比如三个条件同时命中批处理那一列就别再花时间调研实时方案了反过来如果业务明确要求秒级特征那也别纠结先跑个离线版看看效果——直接小步快跑搭实时链路更实际。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的项目就是这么落地的车辆位置事件进KafkaFlink做实时清洗和路径匹配结果同时写Kafka供毫秒级查询和Hudi表供历史回放与分析。历史指标要重算时不再需要一条独立的批处理链路只要重放Hudi里的明细数据就行。坦白讲这种方案在小团队、中等数据体量项目里落地比双链路Lambda舒服得多。5. 三类典型场景的落地拆解从网约车到实时大屏5.1 网约车数据分析哪些环节必须实时哪些环节别碰实时网约车应该是实时离线混合最好的样本了。一个完整的网约车数据体系里数据可以分成好几类选型策略完全不同。订单主状态链路下单、接单、司机到达、行程开始、行程结束、支付——这是典型的必须实时的场景。因为产品端要在用户点开APP的瞬间展示附近车辆实时状态运营端要实时看到各区域订单密度来调度运力。这类数据走KafkaFlink实时清洗写实时宽表或者直接推送到查询服务几乎是刚需。轨迹轨迹和里程数据——这是最有意思的一类。轨迹数据量巨大但上游业务对它的实时要求并不是全局性的司机端行程进行中需要实时位置来做导航和计价预估值但离线分析热力图、常用路线挖掘、里程统计完全可以放到批处理里去算。我在之前的网约车项目里看到过一本经典做法轨迹落Kafka走实时链路给在线地图引擎同时落一份Hive表给凌晨的Spark清洗和分析任务用。财务和计价数据——从我的经验看计价结算这类涉及钱的数据我强烈建议不要为了实时看数把它做成全实时的核心链路。计价本身可以有实时预估给司机端显示一个大概范围但真正的账单生成和对账老老实实走批处理让它可以重跑、可审计。这是我反复跟业务方强调的一条底线实时和钱接触越深出风险的代价越大。5.2 实时大屏与实时公交API秒级更新的伪实时真相很多人一想到数据大屏就觉得背后一定是一套极复杂的实时计算平台。实际上我拆过不少大屏项目——包括用FlaskECharts做的那种经典组合——它们的实时程度往往没你想象得那么高。大屏上展示的订单量、用户量、销售额如果刷新粒度为秒级其实大部分数据并不需要从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%等它变成了业务的一等公民再认真投入不迟。数据架构没有绝对的对错只有在这个时间点、这个资源约束下最合适的选择。
返回列表