ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

实时ETL与批处理ETL选型:从数据延迟到架构落地全解析

实时ETL与批处理ETL选型:从数据延迟到架构落地全解析 干了这么多年数据平台每次有人来问“ETL到底选实时还是批处理”我的第一反应基本都是先反问一句你的业务到底能容忍多少分钟的数据延迟这不是抬杠而是整个选择策略的核心。实时ETL和批处理ETL不是“谁替代谁”的关系它们解决的是不同时效层级的问题。在大数据场景下这两者的工具链、作业形态、容灾方式和运维成本几乎就是两条路选错了后面全是坑甚至会导致项目返工、资源浪费、业务方天天投诉数据不准。这篇内容打算把我在实际项目里踩过的坑、跑过的数据链路、做过的选型评估完整整理出来。无论你是刚接触大数据的学生还是正在设计离线数仓的工程师又或是负责实时计算平台的架构师都可以把它当成一份“决策前的检查清单”来用。我不会只讲概念更多的会落到具体场景、具体参数、具体问题排查上。1. 先把概念盘清楚批处理ETL和实时ETL到底在做什么1.1 批处理ETL的典型工作方式批处理ETL是很早以前就存在的一套数据加工逻辑。它的特点可以用一句话概括定时、批量、集中式处理。你每天凌晨跑一次任务把MySQL里的增量数据读出来清掉无效字段转换格式加载到Hive分区表然后再跑一层Hive SQL去做汇总第二天早上业务方打开报表看到昨天全天的数据这就是最典型的T1链路。这套模式能跑这么多年核心优势在于吞吐能力和可重算性。大批数据一起扫磁盘顺序读、批量写吞吐量可以做得非常高单位数据处理成本也低。最重要的是批处理天然支持“重跑”发现某天数据处理逻辑出错把那个分区删掉再重算一次就行因为上游原始数据还在计算结果也以分区为单位隔离修复成本几乎为零。工具方面大数据场景下通常用Hive做离线的清洗和汇总复杂任务用Spark SQL或者MapReduce调度用Airflow、DolphinScheduler这类框架。Hive本身跑在Hadoop集群上数据落盘HDFS整个链路非常“重”但在大数据场景下胜在稳定、可控。1.2 实时ETL的典型工作方式实时ETL则完全不同。数据一旦产生立刻被采集端发送到Kafka这类消息中间件然后由流处理引擎持续不断地消费做清洗、关联、聚合再写入下游的Redis、Elasticsearch、HBase或者MySQL全程延迟在秒级甚至毫秒级。这里最核心的变化是“作业常驻”。Flink或Spark Streaming启动后就是一个长期运行的进程它不像批处理那样“跑完退出”而是365天挂在那里持续接收数据持续计算结果。也正因为常驻状态管理、checkpoint、容错恢复这些事情就成了避不开的话题。一次网络抖动、一个上游表结构变更、一段反序列化异常都可能造成实时链路数据堆积或结果偏差处理起来比批量麻烦得多。不过实时ETL的价值在于让数据产生价值的时间大幅缩短。在大屏上看到当前实时的交易金额、在风控系统里毫秒级拦截异常行为、在推荐系统里根据最近5分钟的行为调整内容这些都是批处理做不到的。它解决的是“新鲜度”的问题而不是“数据量”的问题。1.3 两种模式的核心差异对照把两种模式摆在一起从几个维度对照能很直观地看到它们各自适合的边界维度批处理ETL实时ETL典型延迟小时级到天级秒级到分钟级处理单位一批数据分区/目录一条条事件流运行方式定时调度跑完结束常驻运行持续消费吞吐量极高适合全量扫描受状态和内存限制适合增量处理成本模型计算峰值集中资源可复用常驻资源占用需要预留Buffer容错重放重跑分区即可依赖Checkpoint和上游Kafka持久化数据准确性容易做到精确幂等重算需要处理乱序、迟到数据语义复杂典型场景离线的经营报表、销售分析、数仓分层实时大屏、监控告警、实时风控、推荐特征这个表格不是绝对的标准答案但可以作为选型前的参考坐标。很多团队一上来就想上实时却发现业务每天看一次数据就够了实时资源完全是浪费反过来有些链路表面看是半小时更新一次业务方却希望延迟越低越好分批跑的话时间卡不准最终还是要上实时。2. 大数据场景下的选型考量是一套多维权衡2.1 时效性需求决定上限先定SLA再讨论架构选型第一步不是选工具而是和业务方把“时效”聊清楚。你要问的不是“要不要实时”而是“从数据发生到用户看到结果最大可以接受多少秒”。这个数字就是SLA。比如经营管理报表前一天晚上生成T1完全没问题反而凌晨跑能避开业务高峰但实时风控如果超过100毫秒响应可能钱已经被盗刷了T1毫无意义。我曾经做过一个网约车数据项目最初产品经理要求“所有指标都要实时大屏”但和司机、乘客的收益结算其实每天只看一次。最后我们把指标拆成两类交易额、订单量、响应时长做了实时大屏所有费用明细、司机收益、里程统计走了T1批处理。这个分类要是弄反了开发成本和资源消耗会直接上涨一个量级。所以我的建议是永远不要基于“感觉”做技术选型而是把业务方要求写进需求文档量化成延迟指标。延迟指标一旦定下来实时还是批量自然浮出水面。没必要让报表展示也走秒级也没必要让风控链路去等着一个跑四个小时的批处理任务。2.2 数据量和成本不是唯一指标但必须算一笔账很多人一提到大数据首先想到的是“量”。量确实重要但成本结构更关键。批处理的成本集中在“跑批那几小时”其余时间资源基本闲置。你可以把这个理解为班车运输固定时间发车一车拉全天的货虽然车大但每天就发几趟资源利用率能被集中调度填满。实时ETL则更像出租车队车辆必须一直停在路边待命随叫随走。即使业务凌晨只有几笔订单作业依然要开着资源依然要占着。Flink集群的TaskManager根据并发度分配内存和CPU Slot只要作业不停止这些资源就一直在消费。数据量每小时几百条和每秒几千条作业并发数可能差不多资源成本几乎一样区别只在每台机器的CPU使用率高低。这是很多项目低估实时成本的地方。做大数据集群部署策略的时候如果链路里既有批处理又有实时任务最好分别评估各自的高峰时段和资源占用。批处理把资源安排在凌晨低谷实时任务要预留高峰期突增的Buffer。预算有限的情况下孰轻孰重一目了然时效要求严格的业务少但关键上实时时效要求宽松的业务多而杂老老实实批处理。2.3 准确性要求和处理语义比延迟更容易被忽略延迟可以靠并发硬扛准确性出问题才是灾难。批处理一般都有明确的分区比如“dt20240115”代表1月15号的数据。这一天处理完数据就不变了后续汇总多用几个分区做关联结果稳定。万一跑挂了你重跑只要源数据不变结果就能完全复现这种准确性保障很简单。实时链路的准确性则麻烦得多。Kafka里的消息可能乱序、可能迟到你早上10点收到的用户点击事件也许实际生成时间是9点59分59秒。你要在聚合时决定“这个事件到底算不算在当前窗口”要配置Watermark来处理晚到数据还要确认下游在故障恢复后不会重复写入。Flink的Checkpoint机制配合Kafka的持久化可以实现端到端Exactly-Once但这是有条件的源头支持消费位点回放、Sink支持幂等写入整个链路要一起配合才有意义。如果业务方说“数据差不多就行差一两条没关系”那你可以大胆选择更轻量、更高效的At-Least-Once模型如果业务方说“最后汇总必须和数据库账一致”那就要考虑用Lambda架构用批处理来兜底纠偏。这里给一个实操判断方法业务方能不能接受“数据先出、后修改”的模式比如大屏先显示99.6%5分钟后修正为99.8%如果可接受实时结果直接用如果完全不能接受结果变化那就不适合纯实时链路。2.4 技术栈和团队能力决定选型能不能落地很多技术选型只盯着“未来的愿景”却忽略了“现在的人”。实时ETL和批处理ETL对团队的要求完全不一样。批处理链路里SQL写得好、会调Spark参数、对数仓分层熟基本就能支撑起来实时链路则需要理解流处理原理能处理背压、检查点、状态后端、窗口语义这些概念还要熟悉Kafka、Flink的监控指标。这不是说实时永远高不可攀而是说项目排期、风险预算都要考虑这个现实。我见过一个数据团队Hive玩得很溜但第一次接Flink光一个“数据倾斜如何定位和解决”就折腾了两周整个实时上线延期了一个月。反过来如果团队里本来就有流处理经验丰富的人即使业务只要小时级更新你完全可以用Flink做微批例如每5分钟触发一次窗口兼顾开发效率和演进空间。真正稳妥的做法是在选型评审时把团队能力写进去必要时先做一个几百行的小Demo让关键开发跑通“Kafka → Flink SQL → 可视化展示”的最小闭环。Demo能在一天内跑通说明技术栈风险可控如果连Demo都糊不上再好的实时架构也只能在PPT里开花。2.5 典型业务场景与建议模式的速查参考多年积攒下来的经验是可以沉淀成一张“选择题答案卡”的业务场景时效需求推荐模式常见工具链路日报、周报、财务对账小时/天级批处理ETLMySQL → DataX → Hive → Spark SQL大屏指标展示秒/分钟级实时ETLCanal → Kafka → Flink → Redis/ES用户行为特征分钟级实时ETLSDK日志 → Kafka → Flink → HBase风控规则命中毫秒/秒级实时ETL埋点事件 → Kafka → Flink CEP → 告警服务网约车订单状态流转秒级纠偏混合架构双链路实时Flink 离线Spark数据仓库基础分层小时级批处理ETLODS → DWD → DWS → ADS这只是一个起点。现实里业务会不断变化今天T1的报表明天可能要求10分钟一更新这就需要在架构上留好演进空间。后面我会专门用一个网约车案例把这种“双链路混合”的落地方式拆开讲。3. 一个具体案例从需求到选择的完整拆解3.1 需求描述网约车订单分析平台网约车数据是典型的大数据综合场景里面既涉及订单的实时状态又涉及司乘的离线结算非常适合用来演示选择逻辑。假设我们正在做这样一个项目平台每天产生几千万条订单轨迹和司机行为日志业务方同时提了几个需求运营大屏上的实时订单量、成交金额、平均接驾时长、区域热力图需要秒级刷新风控模块需要对异常行驶轨迹进行实时识别比如短时间内绕路、频繁取消财务系统每天早上需要昨日全量订单的计费明细、分成明细以及各城市的汇总报表产品部门还要一份历史订单的OD分析起终点分布用于后续优化派单策略跑一次可以一个月不管。如果按最开始那个“非此即彼”的思路要么全上实时要么每天批量跑一次都会有人不满意。全上实时财务对账的结果会因为乱序和迟到数据出现抖动财务那边绝对接受不了全上批处理大屏上的数字基本是“昨天”的运营没法看。这里的选择不是“二选一”而是把一条链路拆成两条按需求分级。实时链路负责秒级指标和风控事件批处理链路负责结算和对账两条链路共享上游数据采集层但各自独立加工、独立存储最后在指标口径上做对齐。3.2 需求拆解后的架构规划整个链路可以简单分成四层数据源层打进数据库的订单表、司机GPS轨迹日志、App埋点日志。采集与传输层数据库变更用Canal监听Binlog写入Kafka日志直接用Filebeat或Logstash采集到Kafka订单数据同时落一份到HDFS用于批处理。实时计算层Flink消费Kafka做实时订单统计、轨迹异常识别写入Redis和Elasticsearch。离线计算层每日定时用Spark读取HDFS上的全量订单清洗后写入Hive分区表再用Spark SQL做财务汇总输出到MySQL。数据源层是公共入口可以避免双链路重复采集。Kafka里既有JSON原始事件也有经过预处理的统一格式HDFS上按“日期小时”路径落一份原始文件批处理直接读实时计算从Kafka消费互不阻塞。这样做还有额外的好处万一实时链路计算逻辑出了问题可以从Kafka或HDFS重新消费修复。3.3 批处理链路细节拆解批处理最重要的是把分层做好。我习惯于把离线数仓按ODS、DWD、DWS、ADS四层来建ODS贴源层原样落地只做数据类型转义和分区裁剪表名建议带ods_前缀。DWD明细层做清洗和维度退化比如把订单状态字段从编码转成中文笛卡尔关联司机基础信息形成宽表。DWS汇总层按城市、司机、小时等维度做轻度汇总产出dws_trip_city_day这类表。ADS应用层跑具体报表指标输出到MySQL或报表系统。一个典型的Spark SQL清洗脚本逻辑是INSERT OVERWRITE TABLE dwd_trip_order_daily PARTITION (dt${dt}) SELECT order_id, driver_id, city_id, CASE WHEN order_status 1 THEN 完成 WHEN order_status 2 THEN 取消 ELSE 进行中 END AS status_name, total_amount, start_time, end_time FROM ods_trip_order_daily WHERE dt ${dt} AND order_id IS NOT NULL;这里有一个批处理特别重要的细节INSERT OVERWRITE TABLE ... PARTITION语法保证了同一分区的幂等性。如果某一天的清洗SQL后续有改动只要重新刷这个分区就行下游汇总表也只重刷受影响的日期分区。不要用INSERT INTO否则重复执行会产出重复数据排查起来极其痛苦。批处理调度上通常用Airflow或DolphinScheduler把任务拆成“ODS导入 → DWD清洗 → DWS汇总 → ADS生成报表”这样一个DAG每个环节依赖上一环节的成功状态。时间线上每天晚上两点拉取昨天的数据三点完成清洗五点完成汇总六点报表数据就能按时上线整个批处理链路的SLA靠调度配置来保障。3.4 实时链路细节实现实时链路的技术栈是 K8s 上的Flink集群 Kafka Redis/ES。用Flink SQL做实时统计比DataStream写Java代码要省很多事我们直接定义Source、Sink和SQL逻辑。订单数据进入Kafka的Topicods_trip_eventFlink作业消费它做两件事第一实时大屏指标。按城市维度每隔60秒计算一次订单量和总金额。这里有一个关键概念是滚动窗口Flink SQL的写法大概是CREATE VIEW trip_5min AS SELECT city_id, COUNT(*) AS order_cnt, SUM(total_amount) AS total_amount, TUMBLE_END(proctime, INTERVAL 60 SECOND) AS win_end FROM trip_source GROUP BY city_id, TUMBLE(proctime, INTERVAL 60 SECOND);因为大屏本身允许几分钟前的数据轻微滞后用Processing Time窗口最简单依赖的机器时钟不用额外处理事件时间排序。但风控的异常轨迹识别就不能用Processing Time了因为迟到事件会漏检。风控链路我们用事件时间Watermark同时加入5秒的允许迟到确保轨迹事件基本都能准确归入对应窗口。第二结果写入。订单统计结果写入Redis通过Flask后端提供API前端ECharts每秒轮询一次刷新大屏。异常轨迹事件写入Elasticsearch供风控平台查询。Kafka的Topic log retention时间设置到了7天为的是实时链路如果出故障可以重置消费位点回放数据重新处理。实时作业的并行度和资源量按高峰期每秒峰值倍数计算。比如订单事件峰值5000条/秒每个并行度处理能力按1000条/秒估算设置并行度为6再预留30%余量。这样既不会因为并行度过大浪费资源也能扛住突发的活动流量。4. 混合架构与常见问题排查实录4.1 为什么我建议先想清楚Lambda而不是非此即彼技术上有个经典争论Lambda架构保留实时批两条链路还是Kappa架构只用流处理用重放代替批量我的看法是对于绝大多数团队Lambda虽然更笨重但活得更踏实。Lambda的实时链路给用户一个“很新但不完全精确”的数批链路给用户一个“精确但稍有延迟”的数两条链路的底层口径统一后最终以批处理结果为准。它确实有开发维护两套逻辑的痛点但它的容错性和调试便利性非常高。Kafka重放、Flink状态恢复都不一定能覆盖所有Bug而批处理重算一个分区就像重新打印一份账本所有人都知道它是最终答案。我见过一个团队刚开始做纯Flink Kappa架构看起来很美后来发现每天要向业务方提供十几个报表这些报表每条链路逻辑都不一样用流计算重写好一遍代价巨大而且一旦数据源有脏数据离线修数比在线修简单太多。最后他们还是不得不引入一版批处理链路。从那时起我对所有实时项目都建议先接受双链路的存在再考虑怎么统一。如果你团队只有实时人才回到Kappa也不是不行但数据对账、补数工具必须提前准备好。4.2 实时链路常见问题与排查技巧实时链路真的是“稳定的时间长了容易忘记它有多脆弱”。以下几个问题我在项目里都遇到过数据倾斜。Flink流里某个城市ID的单量极高比如“北京”这一路Key远超其他城市导致某个并行度处理速率明显低于其他并行度。排查方式很简单看Flink Web UI里每个Subtask的Received Bytes和Busy时间。如果某个Subtask的Busy接近100%大概率就是倾斜。临时方案是给热点Key加随机前缀把数据先分到多个临时Key上再聚合最终再按原始Key进行一次汇总。窗口乱序导致窗口结果不准。比如事件时间窗口是1分钟实际数据迟到了几十秒没有配置Watermark延迟结果少算了一部分订单。排查时看Flink UI中当前Watermark和Event Time的差值修复方式是调整allowedLateness把结果分为“提前结果”和“迟到修正结果”两种下游注意合并。Checkpoint频繁失败或恢复时间过长。常见原因是状态后端存储性能不够或者作业状态太大。建议把RocksDB作为状态后端生产环境配置增量Checkpoint同时把Kafka消费位点存在checkpoint里。如果恢复时间超过预期要考虑是不是单作业状态体积达到了几十GB优化方式是拆分作业按业务域把大的实时链路拆成几个小作业分别管理。背压。某个Sink写ES太慢整个作业积压大量数据。最直接的排查方法是看Flink UI里Source端的背压指标如果一直High就要看下游Sink的写吞吐量。要么并行度增加Sink的并发度要么改批量写入参数要么换一种对写入瓶颈更友好的下游组件。重复消费。实时作业重启后可能从Kafka的最近位点开始消费导致中间一小段数据没处理或者At-Least-Once语义下有重复。解决思路是Sink端做去重比如写Redis用SETNX过期时间写MySQL用唯一键REPLACE INTO这样即使Flink重启多次也不会污染结果。4.3 批处理链路常见问题与排查技巧别以为批处理就永远稳如老狗它的问题无非是显得比较“老”但依然麻烦。分区缺失导致数据质量异常。最常见的就是上游数据迟到了批处理任务当天凌晨启动时只读到了昨天的部分数据整个报表少了一块。排查方式是先去HDFS上看ls里对应分区的文件大小和数量再和前一天对比解决方法是调度上增加上游任务依赖同时设置“等待上游文件就绪”的规则比如每10分钟检查一次分区的数据文件数量和行数达标后才启动ETL。小文件过多导致Spark性能变差。Hive表经常因为上游写太多小文件导致下游Scan慢。我一般会在ODS层用Spark或Hive做一次合并比如INSERT OVERWRITE ... SELECT ...时设置coalesce(50)或者用定期任务把 200MB的文件合并成大文件。小文件不仅影响查询速度还会给NameNode带来内存压力。参数配置不合理导致OOM或GC卡顿。跑大etl任务不要一味用默认配置。Spark SQL中spark.sql.shuffle.partitions建议按数据量设置为200~500之间太大产生大量小map文件太小则单个Shuffle块过大易OOM。还有spark.sql.adaptive.enabledtrue要打开让执行计划动态调整很多查询性能能提升20%~50%。调度依赖没配好导致数据串天。批处理最怕上游还没跑完就开始下游或者日期参数传错了生成一张“必须是昨天但实际是前天”的表。排查方式是在任务日志里记录dt参数和调度触发时间设计时一定要做“上游分区存在性检查”比如Spark作业启动先判断${dt}的ODS分区是否成功写入否则直接失败而不是拿着空数据继续算。4.4 最核心的几条避坑经验最后想着重说几条我在反复踩坑之后才总结出来的经验。第一不要一开始就追求“全链路实时”。先把业务按延迟容忍度分级优先解决核心痛点比如老板要在实时大屏上看订单量那就做这一条实时链路财务对账稳如老狗地跑批就行。全链路实时带来的资源成本和运维复杂度通常是估算值的三倍以上。第二实时和批量必须共用一套数据质量基线。甲乙两条链路算同一个指标结果对不上是大数据团队的噩梦。我建议在Kafka消息里定义统一的事件字段规范批量链路和实时链路都基于这套规范清洗在指标层定期互相对账至少每周比对一次实时汇总值和批处理汇总值的差异范围一旦超过阈值就开查。第三监控体系要尽早建。实时链路除了Flink UI之外还必须有整条链路的端到端延迟指标比如消息从产生到写入Sink的时间还有Kafka消费lag、Redis/ES写入失败率、BackPressure比例。没有监控的实时作业就是盲跑等业务方反馈数据不对的时候可能已经积压了几个小时。第四测试数据一定要包含乱序、迟到、重复、脏字段这些“坏样本”。我见过很多团队造数据全是规规整整的标准字段上线第一天被真实环境毒打。正确做法是在联调环境里模拟一份带异常值的数据流把所有清洗规则和窗口逻辑都验证一遍再复制到生产。我这个人的习惯是真正下决定之前会先在纸上把两条链路画出来标清楚每秒条数、延迟要求、故障恢复时间、可容忍的数据误差。这四列填完选型表格就填完了剩下的只是工程实施问题。如果某个团队还没能把这些数据量化那讨论实时还是批处理就是空谈。希望这篇内容能帮你在下个大数据项目的技术选型现场少一些纠结多一些笃定。
返回列表