ARTICLE DETAIL

资讯详情

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

企业AI平台架构:数据处理全链路设计实战

企业AI平台架构:数据处理全链路设计实战 做企业AI平台的架构设计我从来不在模型层多纠缠真正让我反复返工的永远在数据处理。这句话不是我一个人的感受跟很多团队做技术评审时大家最头疼的往往也是数据接入、特征一致性和训练推理管线这几块。这篇内容就借“AI应用架构师”的视角把企业AI平台架构里的数据处理环节从头到尾拆一遍包括接入层、湖仓底座、特征平台、训练推理管线、框架选型和质量治理。适合正要搭AI平台的技术负责人、数据架构师以及那些模型已经写好、却被数据供给搞得焦头烂额的算法工程师。所谓机器学习中的数据处理本质上不是“把数据清洗干净”而是把原始业务数据加工成可供模型训练和推理使用的特征样本并且这个过程要稳定、可重复、可追溯。很多团队把这块想简单了以为招几个数据工程师跑跑SQL就行。实际做下来你会发现数据处理层的架构决策直接决定模型上线后是稳稳增长还是一路踩坑。1. 先想明白企业AI平台的数据处理到底在解什么题1.1 传统数仓的惯性思维在AI场景下为什么失灵大多数企业已经有一套数仓体系ETL、报表、指标口径都跑得好好的于是做AI平台时第一反应是“在数仓旁边多开几个表”。这个思路不能说错但会漏掉AI场景的三个特殊要求。传统数仓回答的是“过去发生了什么”T1报表、经营分析、财务核算对时效的要求低对口径稳定性的要求高。AI平台回答的是“接下来会发生什么、应该怎么做”它不仅要看历史还要在毫秒级延时内拿到最新数据去做推理。数据形态也不一样数仓输出的是宽表、指标AI平台输出的是样本、特征、embedding这些数据要反复被模型训练任务读取训练时还要保证随机性和分布一致性。另一个被低估的问题是数据血缘的密度。数仓里一张汇总表坏了影响的是一张报表。AI平台里一个特征字段算错了影响的是所有使用这个特征的模型而且模型上线后错误还会被放大。传统数仓的治理经验在这里能复用一部分但远远不够。1.2 三大核心矛盾时效、一致、复用我拆过不少AI平台的需求数据处理层的矛盾说到底集中在三件事上。第一是时效。离线特征按天算在线推理要毫秒级这两者之间必须有一座桥。很多团队一开始只做离线特征等模型要上线时才发现线上没有对应的实时特征服务只能临时抱佛脚。第二是一致。训练时用的特征逻辑和线上serving时的特征逻辑必须完全一致否则模型训练时的AUC再高上线后也会立刻掉点。这个坑我见得太多了后面第五章会详细说。第三是复用。一个企业的推荐、搜索、风控、营销模型可能共用同一批用户基础特征。如果每个算法团队各算各的算力浪费还是小事特征口径打架才是大麻烦——同一个“用户近30天消费金额”两个团队算出来两个数最后业务方都不知道该信谁。理解了这三大矛盾再去看接入层、存储层、特征层、训练管线每一层该解决什么问题就非常清晰了。2. 数据接入层设计让异构系统的数据“进得来、不重不丢”2.1 接入方式怎么选批量、增量还是实时企业里的数据源五花八门业务库MySQL/Oracle、埋点日志、第三方接口、SaaS平台导出的文件、甚至别的部门手工维护的Excel。接入层第一个决策不是选工具而是定策略——哪些数据走批量哪些走增量哪些走实时流。我一般按“业务影响程度”分三档来评估。纯离线分析类的数据比如报表日志、冷启动训练样本T1批量导入就够了没必要为实时付出额外成本。订单状态、库存变化这类数据分钟级准实时增量同步就能满足大多数场景。真正需要毫秒级实时流式处理的通常是风控、反作弊、实时推荐这类对延迟极度敏感的业务。这里有一条我反复强调的原则能批不流能增量不全量。实时流式处理听起来很酷但分布式流系统引入的维护成本、数据乱序问题、状态管理复杂度都是实打实的负担。有些团队一上来就全量接Kafka结果一半topic根本没人在实时消费白白浪费运维精力。2.2 消息管道与断点续传高吞吐背后的可靠性设计确定哪些数据走实时流之后消息管道选型基本就是Kafka的天下了。但很多团队在接入层丢数据问题往往不出在Kafka本身而出在下游消费和写入目标端的方式。Kafka的log机制保证了消息可以回溯消费者通过提交位点offset记录消费进度。理论上只要位点不丢消息就不会丢。实际操作中常见的坑是消费者先写目标表再提交位点目标端写入失败时位点没提交任务重启后会重复消费如果是先提交位点再写目标表写入失败就直接丢数据。无论哪种都说明“至少一次”语义下必须配合幂等写入才能做到不重不丢。我的做法是给目标端统一设计“业务主键数据版本”的去重机制。比如写入Hive表时用主键做upsert写入消息表时以事件ID和事件时间戳做唯一约束。这样一来哪怕是重复消费、任务重跑最终落库的数据也不会重复。2.3 Schema演进与脏数据策略别让上游变更炸掉下游接入层最防不胜防的是上游Schema变更。业务库加个字段、改个字段类型、甚至删个字段都可能让下游解析任务直接报错。我见过最夸张的一次上游DBA把订单表的金额字段从decimal改成string下游特征任务跑了一半全变NULL模型在线效果当晚就崩了。应对方案有两层。第一层是Schema管理统一用Schema Registry管理消息格式变更必须兼容旧版本新增字段给默认值删除字段要提前通知下游。第二层是脏数据分级处理不要碰到解析不了的数据就整个任务失败。可修复的脏数据比如日期格式不对但有规律写进延迟队列做二次处理不可修复的数据写进死信队列并触发告警让值班的人看到。死信队列里积累的数据还能反过来倒推上游数据质量问题比单纯丢日志有用得多。3. 湖仓底座选型训练数据的“粮仓”该怎么搭3.1 湖、仓、湖仓一体从一次架构争论说起存储层是AI平台的地基但很多团队在这里反复摇摆。数据湖便宜、灵活、能存任意格式的原始文件但缺少事务支持和强治理数据仓库治理强、查询性能好但存储和计算成本高模型训练又经常需要直接读原始文件而不是聚合好的宽表。我参加过的一次架构评审会上数仓团队和算法团队吵了一个小时核心分歧就是训练数据到底放湖里还是仓里。最后给出的方案是湖仓一体——一份底层存储既支持BI的数仓查询又允许训练任务直接以文件方式读取原始数据。湖仓一体的核心价值在于“一份数据多种引擎”Spark做批处理、Flink做流计算、Trino做即席查询、训练框架直接读文件大家共享同一份数据不用来回拷贝。这个架构在今天的开源生态里已经非常成熟了新建平台时基本可以闭眼选。3.2 表格式选型Delta Lake、Iceberg还是Hudi湖仓一体的落地依赖表格格式Table Format的选择主流的三个选手是Delta Lake、Iceberg和Hudi。它们都能提供ACID事务、时间旅行、upsert能力但侧重点有差异。我列过一张对比表方便你在选型时快速对齐维度Delta LakeApache IcebergApache Hudi事务与快照强基于Spark生态强多引擎兼容好强对写入优化深入upsert能力支持merge语法成熟支持但部分引擎需适配支持MOR表实时读优化多引擎支持主要是Spark其他靠适配Flink/Trino/Spark兼容性最好Spark/Flink都支持CDC入湖场景一般一般场景设计最贴合运维上手成本较低文档多中等中等偏高如果团队以Spark为主、希望快速落地Delta Lake的上手曲线最友好。如果数据要同时给Flink、Trino、Spark多个引擎用Iceberg更中立。如果业务有大量CDC同步、需要低延迟的增量读取Hudi的MOR表会更顺手。没有绝对银弹关键是看你现有引擎生态和团队运维熟悉度别为了追新而追新。3.3 分区布局与小文件治理决定读写效率的隐藏细节存储层最容易出问题的是分区策略和小文件失控。分区字段选错扫描数据量能差一个数量级。比如订单表按天分区特征任务却经常要扫最近90天那就应该再设计一个按业务维度如用户ID范围的二级分区策略避免每次全表扫描。小文件问题在AI场景尤其突出。训练数据动不动几亿行但很多任务写出来的是几MB甚至几百KB的小文件Spark读文件的开销比计算本身还大。我之前接过一个特征回填任务原始数据200GB小文件两万多个跑一次要4小时后来做了一次compaction合并成上千个大文件直接降到40分钟。这类问题没有一劳永逸的解法要靠设计层面去控制写入任务合理设置并行度避免过度分区定时跑compaction作业合并小文件增量任务尽量复用同一个分区目录而不是每次新建。存储层的健康程度决定了整个AI平台数据处理管线的下限。4. 特征平台把“喂给模型的数据”从工程资产变成产品4.1 为什么特征层要单独拎出来建设模型训练需要的不是“数据”而是“特征”。同一个用户ID下的近30天消费均值、最近一次登录距今天数、商品embedding的相似度这些都是特征。传统数仓不会沉淀这些口径于是每个算法团队都从原始表开始自己算重复开发不说口径还经常对不上。特征平台的核心价值有两个一是特征复用一个特征定义好之后任何模型、任何团队都能通过统一接口获取二是口径统一同一个特征在任何模型里的数值必须完全一致。这听起来像是工程洁癖但实际影响非常大。我见过有的公司两个模型都用了“用户活跃度”这个特征一个按周定义、一个按月定义最后两个模型在同一批用户上的预测结果完全无法横向比较。特征平台建设不一定非要买商业产品开源的Feast、阿里的JindoFS方案、或者基于FlinkRedis自研都能达到效果。关键是想清楚边界哪些特征进平台、谁负责维护、变更流程怎么走、离线在线怎么对齐。这四个问题不定清楚平台搭起来也是摆设。4.2 离线特征与在线特征的一致性一套逻辑两个世界这是特征平台最核心、也最容易翻车的环节。离线训练时用Pandas或Spark算特征线上serving为了性能又用Java或C重写一遍两边逻辑稍微不一致模型上线就会出问题。我处理过的一个经典案例一个推荐模型的“用户近7天点击次数”特征离线用Spark窗口函数算线上为了降延迟用Redis实时累加结果线上数值总是比离线小一点。查了半天发现是时区问题——离线按自然日边界切分线上按UTC切分两边差了8个小时。这种问题在训练时完全看不出来上线后点击率掉了3个点才定位到。现在的标准做法是“一套逻辑两个世界”特征计算的核心逻辑用跨平台语言统一实现比如Flink SQL或共享的特征计算库离线批量计算和在线实时计算复用同一份代码再配合离线在线特征对比工具抽样一批样本验证两边计算结果是否一致误差在阈值内才允许放量上线。4.3 实时特征与特征回填窗口、乱序和重算的艺术实时特征和离线特征有个本质区别离线计算面对的是完整数据实时计算面对的是不断到达的流数据。常见窗口聚合场景比如“最近5分钟点击次数”就绕不开乱序和迟到数据问题。Flink的watermark机制就是为这个设计的。你需要根据业务容忍度设置水位线和允许延迟时间比如点击日志允许迟到30秒超过30秒的数据默认丢弃或走侧输出流单独处理。这里没有无脑配置的答案延迟容忍度和特征准确性是直接权衡风控场景可能只容忍几秒营销分析就能放宽到几分钟。特征回填则解决另一个问题当特征口径变更或者模型要重新训练时需要重新计算历史一段时间的特征。回填任务必须支持按天、按小时粒度的重跑而且每次重跑要幂等——同一个时间分区不管跑多少次结果都一致。我在回填任务里通常加上“输出分区先清空再写入”的逻辑宁可慢一点也要保证数据干净。5. 训练到推理的数据管线一致性是最大的隐形坑5.1 数据集版本化模型可复现的底层保障算法工程师经常说“我这个模型效果不错”但问一句“训练数据是哪份”就卡住了。很多团队训练数据就是“昨天从这个库里捞的”没有版本、没有快照、没有记录。等想复现一个效果时原来的数据已经变了模型结果对不上排查无从下手。我推荐的做法是给每次训练建立一个数据集描述文件dataset manifest用JSON记录数据源表、时间范围、特征版本、采样方式、清洗规则、过滤条件。下面是一个最小示例{ dataset_id: rec_v1_train_20240601, source_tables: [dwd_user_action_di, dwd_order_di], time_range: [2024-01-01, 2024-05-31], feature_version: user_features_v3, sampling: {method: negative_downsample, ratio: 0.1}, filter: user_id % 10 0, created_by: algo-zhang, created_at: 2024-06-01T10:00:00Z }每次训练开始前先记录manifest对应的训练数据目录设置成不可变快照。这样任何一次训练结果都能追溯跑完发现效果差先看数据版本对不对换特征了对比两个版本的特征差异。这层能力建设成本不高但价值极大。5.2 训练与推理数据偏斜线上线下的“两张脸”训练与推理数据偏斜Train-Serve Skew是AI平台数据处理里最隐蔽、最头疼的问题。训练时看到的用户分布、特征分布和线上推理时完全不一样模型自然就偏了。最常见的原因有几种线上某些特征取不到值用了默认值填充但训练时这些字段根本没处理过空值线上特征实时计算时数据还没到达比如用户刚产生的行为实时特征窗口根本覆盖不到还有代码版本不一致线上跑的是一套老逻辑训练用的是新逻辑。我的排查经验是在模型上线前用最近一段时间的线上真实请求日志做一次回放把每条请求输入到特征服务生成一份“线上特征样本”再和离线训练样本做对比重点看特征缺失率、默认值占比、分布偏移PSI。只要这一步过了大部分Train-Serve Skew问题都能在上线前暴露出来而不是等线上出事故再查。5.3 数据回放与评测集管理用历史请求替模型“体检”评测集管理看起来简单但很多团队做得草率。评测集不能只取最近一个月的数据因为业务有周期性大促、淡旺季、节假日的数据分布差异极大。我的建议是评测集要刻意覆盖多个业务周期正常日、大促日、极端天气日、甚至历史上出现故障的异常日都放进去这样模型评估才不是粉饰太平。数据回放是另一个强有力的验证手段。把历史上真实的请求日志包含当时的上下文特征重新灌入新的特征服务或模型服务观察预测结果是否符合预期。回放的价值在于不需要等流量灰度就能在测试环境里提前发现特征服务的时序错误、数据缺失、版本兼容问题。这套机制建立起来后模型上线前体检就是固定动作数据回放、特征对比、评测集跑分三项全部通过才允许发布。坚持一段时间后你会发现线上事故率会明显下降模型回滚次数大幅减少。6. 数据处理框架选型与实战调优Spark、Flink、Ray怎么选6.1 三个框架的适用边界数据处理框架的选型经常被搞成“信仰之争”但其实各自边界很清晰。Spark适合大规模批处理、离线ETL、批量特征计算生态成熟、稳定性好AI平台里80%的离线任务都应该跑在Spark上。Flink适合实时流式处理、窗口聚合、在线特征计算这是它的主场。Ray则更偏训练侧适合RL、AutoML、分布式数据预处理这类需要灵活调度的场景和模型训练框架配合更紧密。我做得一张选型参考表框架核心模型语言生态典型场景不适合的场景Spark批处理Scala/Python/SQL离线ETL、批量特征、大规模join毫秒级流处理Flink流处理Java/Python/SQL实时特征、窗口聚合、CDC管道大规模批处理Ray分布式调度Python训练数据预处理、RL、AutoML标准SQL ETL选型时不看框架名气而是看你的任务属于哪一类。把Spark的任务硬套到Flink上或者用Ray去跑标准SQL ETL都是给自己找麻烦。6.2 作业调优实践数据倾斜、分区与资源参数框架选完真正的硬仗在调优。我见过太多团队遇到作业慢就加资源加完还是慢最后发现是数据倾斜或者分区爆炸。数据倾斜是Spark作业最常见的问题。一个热点key占了90%的数据单个task跑完要一小时其他task几分钟就结束了。处理方式有几种给join的key加随机盐salt、广播小表、对热点key单独拆分处理。我在一个用户行为特征任务里碰到过“某头部用户一天产生百万级行为数据”的情况对热点用户单独拉出来计算再合并整体作业时间从2小时降到30分钟。分区和内存参数也不要盲目拍脑袋。通用经验是shuffle分区数设为executor总核心数的2-3倍内存配比按“executor内存:shuffle内存:rdd内存”的经验比例设置而不是直接扔一个spark.executor.memory20g。Flink那边重点看checkpoint间隔和状态后端大状态场景选RocksDB小状态场景用堆内存就够了。调优的核心是先看数据量和数据分布再定参数别默认“加机器万能”。6.3 高通量场景的稳定性三件套重试、幂等、限流高通量数据处理场景下稳定性比性能更重要。日增几十亿条消息的管道任何一次抖动都可能造成数据积压或丢失所以必须默认“故障会发生”来设计。重试要有指数退避策略不能一失败就立刻重试那样只会打爆已经脆弱的系统。幂等是重试的前提前面接入层说过目标端要有业务主键去重否则重试一次就重复一份数据。限流和熔断则是为了保护下游上游消费速率突然翻倍时不能把压力全部怼到存储层或在线服务上要在管道中预设阈值超过就降级。我记得一次大促前做全链路压测一个实时特征管道因为下游Redis抖动所有写入任务都在疯狂重试反而把Redis彻底压垮了。后来加了熔断器下游异常时自动切换为读降级缓存等下游恢复再补写大促当天数据管道稳如老狗。高吞吐场景的稳定不是靠某一次修bug而是靠这些机制组合兜底。7. 数据治理与质量门禁容易被低估的隐形护城河7.1 质量门禁在数据进入特征和样本之前把住关口很多团队把数据质量当成事后补救出了问题再捞日志分析。但AI平台的数据质量应该是事前防御——在数据进入特征计算和样本生成之前先过一道质量门禁。我通常设置四类检查指标字段完整性非空率低于阈值就阻断、主键唯一性重复率异常就告警、分布稳定性计算PSI超过0.2就触发检查、跨表join命中率关联主表掉数据会严重影响样本质量。每一项检查都有对应的阻断和告警动作可自动修复的走修复流程无法修复的直接阻断下游任务并通知责任人。一个具体的例子某次上游埋点改版用户行为日志的“页面ID”字段有30%变成了空值。放在以前这会导致接下来的特征任务、训练任务全部白跑第二天才会被发现。有了质量门禁后凌晨1点的检查作业发现非空率跌破阈值自动阻断下游并告警算法团队上班时已经在排查问题了。这套机制看起来朴素但价值极大。7.2 血缘与可观测性出问题时你能否快速定位数据平台做大了之后最痛苦的不是没有能力解决问题而是不知道问题出在哪一层。一条数据从业务库到训练样本中间可能经过接入、清洗、特征、样本生成五六道工序。某个模型效果异常到底是哪一环出了问题血缘系统要回答两个问题这条数据从哪来、用到哪里去。正向看能追溯样本里每个字段的加工链路反向看能知道某个特征表被哪些模型引用、下线一张表会影响多少模型。平台有了这个能力任何一次变更的影响范围都能提前评估。可观测性方面我给每个处理任务都埋了基础指标输入行数、输出行数、耗时、失败率、脏数据量。这些指标不是为了报表好看而是为了排查问题时能一层层下钻先看管道层有没有积压再看质量门禁有没有阻断最后看具体任务日志。三层定位下来大部分问题都能在十分钟内找到根因而不是开全员大会互相甩锅。7.3 数据安全与脱敏接入层就要做的技术底座最后说数据安全。AI平台处理的数据往往比传统数仓更敏感因为它包含用户行为、交易记录、位置轨迹等个人信息。这类数据的脱敏和权限管控必须作为架构的一部分在接入层完成而不是等数据到了训练环节才想起来。技术层面有几件基础工作敏感字段统一做脱敏处理常见方式有哈希脱敏、令牌化、可逆加密列级权限按角色严格控制算法工程师默认拿不到手机号、身份证号等明文训练样本导出时自动过滤敏感字段日志文件不得打印原始个人信息。模型推理时如果需要使用某些敏感特征也要通过特征服务统一处理而不是让算法同学自己拼SQL去取数。把数据安全前置到接入层后面所有环节都不用再担心“这数据能不能用”的问题。我也见过一些团队在数据平台搭得差不多了才发现脱敏没做结果被迫返工——这种返工成本极高而且容易返出新的数据质量问题。做AI平台架构这些年我最大的体会是不要一上来就追求全链路实时也不要堆一堆炫酷的组件。把离线批处理做稳、把特征口径统一、把质量门禁装上、把训练推理一致性管住这几件事做扎实了比引入十个新框架都管用。数据处理这一层看起来不性感但它离业务结果最近是整个AI平台最不该偷工减料的地方。希望这篇梳理能帮正在设计数据平台的你少走几步弯路。
返回列表