ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

从传统ETL到全域数据平台:架构演进、实时集成与湖仓一体实践

从传统ETL到全域数据平台:架构演进、实时集成与湖仓一体实践 做了十几年数据集成我越来越觉得ETL这个词已经装不下现代数据平台的复杂度。早年我们谈ETL就是抽取、转换、加载三个动词撑起一整套数仓。现在再聊数据集成架构你得面对实时流、湖仓一体、数据血缘、数据服务化、元数据治理这些东西甚至还要跟上AI辅助建模的节奏。这篇文章想把我自己在一家零售企业做数据平台架构演进的过程做个复盘从老旧的存储过程式ETL一步步迁移到今天的全域数据平台。全程没有堆概念都是踩过坑之后的真实取舍。1. 从一个真实改造项目说起ETL的日子为什么难熬1.1 那个每天凌晨两点才跑完的批处理任务我当时接手的是一个典型零售集团的数据平台底层有SAP ERP、门店POS、会员CRM、线上商城、第三方支付渠道零零总总十多个业务系统。传统ETL架构用的是SQL Server 存储过程每天晚上从各业务库抽取增量数据经过清洗、关联、汇总写进数仓的明细层和汇总层。听起来很常规但真实情况是这套批处理链路有四十多个数据同步任务任务之间存在依赖关系比如订单明细没跑完日交易汇总就不能启动。最惨的时候某个上游业务库凌晨一个大查询把资源占满同步任务直接超时后面整个链路全部卡住。值班同事凌晨两点被电话叫起来手动补数、重跑依赖任务弄到天亮才把昨天报表的数据凑齐。这种日子你问做数据运维的人谁都不想再来一遍。更折磨人的是需求响应。业务要一个新指标比如门店周边三公里会员的复购率听起来不就是拉两个表join一下但实际要追溯会员地址清洗逻辑、门店分区归属表、复购周期口径定义少则半天多则两天因为数据链路不透明谁都不敢拍胸脯说某个字段的口径是对的。1.2 传统ETL的三个致命伤做了几年之后我把这类型传统架构的问题归纳为三个致命伤后面所有改造动作几乎都是围绕这三个问题展开的。第一是存储计算强耦合。那时候数仓跑批用的机器和数据库是绑在一起的存储过程直接在数据库引擎里跑。业务量一涨计算慢了唯一能做的就是换更强的机器横向扩容几乎等于重新迁移一遍。而到了全域数据平台阶段对象存储负责存、弹性计算集群负责算两件事完全分开资源不够了直接加计算节点就行。第二是实时性严重不足。传统ETL默认是T1当天数据凌晨出报表。但现在的业务决策比如线上促销活动效果、异常订单拦截、会员积分异动提醒T1根本来不及。你凌晨告诉运营昨天转化率跌了活动已经结束了有什么用所以后来我们做全域平台第一条硬指标就是关键数据从业务事件发生到可查询不超过5分钟。第三是数据标准系统性缺失。一个订单状态字段在ERP里是数字编码在报表里是中文文案在CRM里又是一套枚举值。传统ETL靠存储过程在中间转换但这些转换逻辑散落在上百个脚本里没人能说清全局是什么样。这个痛点最终倒逼我们把元数据管理和数据建模规范提升到架构级而不是停留在脚本注释里。2. 全域数据平台到底全在哪逻辑架构与关键能力拆解全域数据平台这个词这几年被用滥了很多厂商把数据中台换了个皮就说是全域。我这里讲的不是营销概念而是工程上真实跑通的逻辑架构。整个平台从下往上分五层数据接入层、湖仓存储层、计算处理层、数据服务层、治理与安全层。每一层都有和传统ETL完全不同的设计逻辑。2.1 数据接入层从轮询抽取到事件驱动传统ETL的数据接入方式是按固定时间轮询比如每小时去数据库查一次增量字段或者干脆每天全量拉。这种模式最大的问题是不知何时发生变化你得靠频繁查询去弥补盲区既浪费资源又容易错过数据。现代全域平台的数据接入层核心设计原则是事件驱动。业务数据库产生一条变更记录通过CDCChange Data Capture变更数据捕获机制把这条变更同步到消息队列下游所有需要数据的系统都能实时感知。例如订单被取消这一条事件会立即触发库存回补、渠道对账、用户积分冲正等多个链路不再需要等一个批处理任务统一处理。我们在实际改造中对比过两种方案的投入产出。轮询方案看起来简单但每次查询都会对业务库造成压力到了大促高峰期DBA总会来投诉说集成平台把核心表锁住了。换成CDC之后源库只需要开启binlog同步工具订阅日志即可生产库的额外负载可以忽略不计。这个转变是整个架构从被动等数据到主动流数据的关键一步。2.2 湖仓一体存储不再纠结先有库还是先有仓过去数据架构里数据仓库和数据湖是两套东西数仓负责存储结构化的、用于报表和BI的数据数据湖负责存原始日志、图片、文本等非结构化数据。两套存储之间还要做一层同步维护成本极高而且经常出现数仓里拿不到湖里的数据湖里的数据不敢直接用来出报表的尴尬。湖仓一体解决的就是这个割裂问题。它保留数据湖的低成本存储和灵活格式又引入数仓的ACID事务、Schema约束和高效查询能力底层表格式比如Iceberg、Hudi、Paimon直接承担了统一元数据的工作。现在我们在平台上既能把几十TB的点击流日志原样落进湖中又能在同一个表上跑SQL做精确对账不需要再做第二份拷贝。从架构演进角度看不要再纠结上数仓还是上数据湖——它们是同一个底座的不同产品形态。真正要规划的是原始区域、明细标准区域、汇总服务区域如何分层访问权限如何隔离。这个思路在改造的早期定下来后面所有业务接入都顺着这套分层走很少出现推倒重来的情况。2.3 数据服务层让下游拿数像点外卖传统数仓的消费方式是开权限连数据库自己查——一条SQL写不对能跑个小时还把集群拖垮。全域平台一定要把这层做成数据服务化下层数据通过统一接口暴露上层应用通过API或者指标平台获取数据。我们平台里的数据服务层是指标API双通道。高频使用的交易额、客单价、库存周转天数这些指标全部在指标管理平台里统一定义业务人员直接拖拽查询应用系统需要明细数据则通过数据API授权后按参数取数。这样做不仅隔离了下游对底层物理表的依赖还能精准控制查询的资源消耗额度个别慢查询被熔断也不会影响整个平台。有个数据服务化的天然好处是口径收敛。以前十张报表有十个销售金额因为各自join的粒度不一样。现在指标平台只认一套逻辑下游想改口径必须走审批数据可信度大幅提升。这一层是整个平台建设里最容易被忽视的但它恰恰决定了平台价值能不能被业务感知。2.4 元数据与数据血缘平台的可观测基础传统ETL时代跑批失败看重新跑就好但更隐蔽的问题是这个字段从哪来、那个指标怎么算。没有元数据体系和数据血缘整个平台就是一个黑盒出了问题只能靠人肉翻脚本。现代全域平台的元数据管理分两类技术元数据和业务元数据。技术元数据包括表结构、分区信息、同步任务状态业务元数据包括指标口径、维值含义、数据所有者。数据血缘把这两类连接起来一张报表里的销售额可以从物理字段一路追踪到业务系统的订单表。我们落地血缘后的一个典型场景上游业务系统做了一次字段调整平台通过血缘分析自动提示以下17个下游表和6个指标会受影响。这在传统架构里几乎不可能实现得靠老员工回忆这个表谁在用。数据血缘加上自动化影响分析是从头痛医头式运维变成平台级可观测的核心能力。3. 演进路径上的关键决策技术选型与架构取舍从老ETL往全域平台迁最大的困难不是技术不会用而是每一层都有好几种选择选错一条后面都是坑。我把自己做过的选型思考列出来偏向中小规模的团队但思路是通用的。3.1 自研还是买商业套件常见的第一道选择题。自研数据集成框架灵活度高能深度适配内部流程但成本巨大——需要维护连接器、调度、监控、权限、版本升级一套下来没有三五个人持续投入根本跑不稳。买商业套件省事但容易被厂商锁定定制改造要排期而且商业套件对新兴数据源的支持往往跟不上。我个人的建议是混合策略核心调度与编排自研因为这是我们独有链路逻辑、改造最频繁的部分通用连接器和底层存储直接用成熟开源组件没必要重复造轮子。实际我们用的调度引擎就是基于开源项目二次开发的只保留了自己需要的节点类型和告警触发逻辑半年的时间投入换来了完全符合业务场景的编排能力。一些团队喜欢全自研全家桶最后发现光适配各种数据源格式就忙不过来。另一些团队什么都买商业套件大促前需要加个新数据源结果得等两周排期。混合策略在速度和掌控力之间相对平衡。3.2 实时与批处理的融合Lambda、Kappa还是流批一体这是架构演进里绕不开的流派之争。Lambda架构维护两套代码——实时链路跑流处理离线链路跑批处理保证最终一致但开发运维成本高口径还容易漂移。Kappa架构用实时流引擎同时处理实时和批理念上只维护一套代码但要支撑大规模回放时对消息队列的存储能力要求极高不是所有团队都扛得住。我们最终的落点是流批一体底层存储用Iceberg统一表格式计算引擎既能跑批作业也能跑流作业。这样一套表结构、一套逻辑既能满足T1的历史全量统计又能支撑十分钟粒度的实时指标。流批一体的前提是表格式要支持增量提交和过期快照这些Iceberg都天然具备。从工程实施角度不要一上来就追求全实时那是自找麻烦。最稳妥的路径是把业务按实时刚需程度分级真正的实时需求促销监控、风控、库存联动走流处理绝大多数分析报表保持微批比如五分钟一次历史回刷这类场景走离线全量。同一个平台同时支持这几种节奏比逼自己一条流打天下舒服得多。3.3 常见组件选型对照我整理了一份我们在设计平台时反复比较过的核心组件对照不是评测标准答案但基本代表了当前主流的架构取向功能域候选组件我们的选择与理由数据接入/CDCCanal、Debezium、Flink CDCFlink CDC为主因为它能同时支持全量快照和增量日志还能配合后续流任务跑数据转换消息队列Kafka、Pulsar、RocketMQKafka生态最成熟Flink整合也最顺手Pulsar更灵活但如果团队不熟悉运维成本高实时计算Flink、Spark Streaming、StormFlink状态管理、窗口机制、checkpoint能力更强做流批一体的基础批量计算Spark、Hive、Presto/TrinoSpark做离线重处理Trino做交互式查询两者互补湖仓表格式Iceberg、Hudi、PaimonIceberg对传统数仓人员更友好快照隔离和Time Travel查询真的省心即席查询引擎ClickHouse、StarRocks、DorisDoris统一了明细模型和聚合模型还能直接兼职做数据服务层的加速查询调度编排Airflow、DolphinScheduler、内部引擎内部引擎基于DolphinScheduler二次开发更贴合实时任务和批任务的统一DAG数据服务网关APISIX、Kong、自研薄封装APISIX稳定性好网关层面的限流、鉴权都能覆盖关于选型有一个重要经验不能只看性能测评数字还要看团队熟悉度和运维体系。我们最早想上Pulsar因为它的多租户能力很强但运维团队没人扎实踩过Pulsar的坑最终还是切回Kafka把精力放到更好的监控和扩容策略上反而更稳。3.4 计算与存储分离的底层逻辑全域平台与传统数仓最大的物理架构差异是计算与存储分离。传统数仓一台机器又存又算数据量上来后很难平衡。新平台用对象存储比如MinIO或云上对象存储作为统一底座计算引擎通过连接器按需读取数据完全不依赖本地盘。这个架构对数据集成场景的直接影响体现在扩容与缩容上。比如大促前一天需要跑全量用户画像我直接临时扩充200个计算节点大促结束后缩回20个存储不受影响。传统ETL想这么干几乎不可能因为你没法把一个SQL Server集群在一天内弹性扩容再在三天内缩回去。存储与计算分离也带来了数据访问模式的改变外部系统可以直接从存储层获取数据文件而不是非要通过计算引擎。比如数据分析师要用Python读一个CSV格式导出的数据集平台可以直接授权他访问对象存储的特定目录连SQL都不需要。集成平台的数据分发能力就是这样从表对表变成目录对表、桶对任务的弹性格局。4. 数据建模与质量治理平台不是搬数据那么简单很多团队做平台迁移时天真地以为把同步任务从存储过程改成Flink作业、把表从SQL Server搬到对象存储就是演进了。结果数据搬过去之后报表还是对不上口径还是各说各话甚至因为实时链路的引入错误数据扩散得比以前更快。数据建模和治理跟不上架构层面再先进也是虚的。4.1 为什么传统数仓建模在实时场景会失效传统维度建模星型模型/雪花模型在T1场景下很好用因为你有大把时间在批处理里做清洗、维度退化、缓慢渐变维处理。但到了实时场景数据是连续流进来的你不可能在一条流里同时维护几十个维度表的缓慢变化历史那会让状态无限膨胀。我们踩过的坑是一开始照搬维度建模方式做实时事实表结果订单数据流一进来就想实时去关联会员维度和商品维度Flink的维表join先是查MySQL后来改成查Redis再后来发现维度需要支持历史版本不得不用HBase存储维表快照。链路越长实时性越差而且维度值一变历史事实全部跟着变——这在报表里是灾难。后来我们调整为宽表优先、维度后补实时链路只保留最核心的原始事实字段维度关联全部推迟到服务层做。也就是说实时明细表里存的是会员ID和商品ID等业务方查询时API网关再通过维度服务去补充名称和分类。看起来不优雅但实际效果最好既保证实时事实不膨胀又让维度信息永远是最新的。4.2 质量规则前置入口校验与在线监控传统ETL的数据质量通常靠事后抽检发现脏数据时已经进了报表再跑一次修复脚本整个过程又慢又不透明。全域平台必须把质量校验前移数据进入消息队列之前做基础校验进入湖仓存储时做规则校验生成指标之前做业务合理性校验层层设卡。我习惯把质量规则分成三档。第一档是物理规则字段非空、主键唯一、数值范围、长度限制这些在写入存储之前必须全部通过。第二档是业务规则比如订单金额不能为负数、会员年龄区间是否合理这类规则直接影响指标可信度。第三档是统计规则和数据分布有关比如某支付渠道的订单量突降50%系统自动预警而不是等运营发现数据不对。平台上线后我做过一次统计90%的数据异常在入口校验阶段就被拦截真正需要事后补救的只有少数慢变数据。这比传统跑批之后看日志的效率提升是数量级的。质量规则一定要配上自动预警和告警分派否则规则堆积再多没人看也是白搭。4.3 数据目录和安全管控的落地细节全域数据平台因为数据集中了安全责任也集中了。早期我们的教训是权限粒度太粗一个数据管理员角色能读全库的明细数据非常危险。现代数据平台要求做到字段级、甚至值级的数据权限控制。我们在平台上落地了数据目录分级分类按需授权的三层结构。数据目录回答平台有什么数据每张表都标注负责人、业务域、敏感等级分级分类把数据分成公开、内部、敏感、机密四类申请权限必须说明用途和使用期限平台自动审批流程走完只开放必要字段和行级范围。脱敏也是一个容易被忽视的细节。实时流和离线表的数据脱敏必须保持一致否则实时接口返回明文离线报表却是脱敏后的值业务侧会疯掉。我们统一把所有敏感字段的脱敏算法放到底层视图层执行不管数据走实时API还是离线分析最终面对用户的结果都是同一套脱敏规则。这个细节看似工程小问题却直接影响数据共享的顺畅度当时调试了整整一周。5. 分阶段迁移的实操路线从试点到全量替换架构设计得再漂亮迁移时太激进一样会翻车。我们当时定的原则是链路优先、平台并行、指标验收、逐批切换整个过程大约用了六个月经历了无数个重新对数的夜晚。这段实操经历我认为比任何架构图都更有参考价值。5.1 第一步以取数链路为单位做梳理不要按表来迁移因为很多表之间是强依赖的单独搬一张表没有意义而按链路迁移可以保证端到端验证。具体做法是把业务数据流拆成若干条主干链路比如订单履约链路会员生命周期链路商品库存链路。每条链路梳理出源系统有哪些表、中间做了哪些转换、下游产出哪些报表和接口。我们为此建了一张链路梳理的共享文档每一节都标注当前新旧系统的状态、负责人、已知的口径差异。这个过程看起来像考古学但确是后面一切工作的底座。整理链路时你会发现很多历史遗留的不知道谁在用的表别急着清理先标记待确认。等到确认不需要了再下线避免误伤业务。我们最终清理了约15%的废弃任务为平台减负不少。5.2 第二步双跑与对账机制的搭建迁移期间最安全的策略是新老平台并行跑数据每项关键指标两边口径输出必须严格一致。如果对不上那就先在老平台查原因确认是新平台逻辑问题还是老平台本身口径就错了。对账机制要覆盖三个层次行数对账、主键对账、指标对账。行数对账能发现数据丢没丢主键对账能发现重复和更新异常指标对账能发现计算逻辑差异。我们最开始对账是写SQL两表countgroup by效率低下后来开发了统一对账任务每天自动比对关键表的count、sum、distinct key值有差异直接发钉钉告警。这里有个容易忽略的坑时区问题。新平台接入的消息事件是UTC时间戳老平台直接拿本地时间存进库两边对账时同一个订单的下单时间差8小时导致日分区归属对不上。这个坑排查了整整两天最后才发现是时区转换遗漏。建议所有平台统一使用UTC存储、展示时再转本地时区这是集成架构里最基础但最容易犯的错。5.3 第三步切换路上的典型踩坑切换阶段我们遇到了很多正常文档里不会写的问题挑几个特别典型的说说。第一个是任务依赖丢失。老平台的调度器有隐性依赖比如A任务今天跑完B任务必须在凌晨1点前完成但这个依赖关系只存在于人的脑子里新平台的DAG并未配置。有一次切换后某个报表延迟了三个小时因为调度器认为B任务没有依赖、可以立刻跑但它实际需要等A任务的数据。解决方法是把旧平台所有任务依赖手动梳理进新调度器并在切换初期保留告警阈值比平时放宽一点避免误报。第二个是重复数据处理。实时链路天然容易产生重复消息比如源库网络闪断后重发binlog消费端没做幂等就会出现同一笔订单被记两次。新平台要求每张事实表必须建立业务主键的去重机制比如通过Flink的rocksDB状态做去重或者依赖Iceberg的row-level操作。别指望消息队列恰好一次就能保证语义消费端保护必须做足。第三个是文档和实际不一致。老数据仓库的注释时间久了就失真比如一个created_at字段实际上存的是更新时间。迁移时按注释取数结果对账对不上最后查源码才发现字段语义变了。这里没太好的捷径只能靠双跑对账发现问题再倒逼业务确认字段真实含义。5.4 复盘平台上线后我做的三件事平台稳定运行之后我没有急着把剩下的任务全部切完而是先做了三件事每一件都在后续持续产出价值。第一建立了一套数据集成任务的黄金指标看板同步延迟、消息积压量、脏数据拦截率、任务失败恢复时长。集成平台原来是被动响应数据怎么没更新现在这些指标一出来主动解决问题比业务发现快了半天。第二把之前踩过的坑写进了团队的《数据平台开发手册》。包括时区处理规范、幂等设计模板、血缘标注规范、权限申请模板。新同事上手时先读这份手册明显减少了很多低级错误。第三推动业务侧养成了指标评审习惯。任何新指标上线前先找数据负责人确认口径、来源和血缘再进平台开发。这个过程虽然要开会但比上线后数据对不上要省时间得多。回看整个演进过程技术选型从ETL升级为流批一体平台只是表象真正深刻的变化是思维方式的改变数据不再是每天定时搬运一次的商品而是从业务事件发生那一刻就需要被感知、被管理、被服务的资产。这个认知一旦确立架构演进自然就发生了。最后再分享一个实操建议如果你的组织还停留在传统ETL阶段别急着全盘否定它先把它最痛的三个场景挑出来用新的集成架构做两个试点让业务看到实时数据的价值再推动全量演进阻力会小很多。
返回列表