十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

AI数据中枢实战:从架构设计到数据协同落地

AI数据中枢实战:从架构设计到数据协同落地 我前阵子带团队做完了一个内部代号叫“统好AI数据中枢”的项目核心就一句话把散落在各业务系统里的数据统一收口让它们能相互协同然后喂给AI应用和大模型驱动整个企业往数智化方向走。这类项目在业内已经不算新鲜但真正把“数据协同”这四个字落到实处的很少大多数还是停留在搭了一个数仓、接了十几张表、跑几个BI报表的阶段。这篇就围绕这个项目讲讲我是怎么拆解的、架构怎么设计、协同引擎怎么落地、哪些坑是我实际踩过的希望给正在做AI数据基础设施的同行一些参考。项目立项时我们面对的真实情况是企业已经有ERP、CRM、MES、HR系统数据量不小但每条业务线各管各的财务报表口径和运营报表口径对不上AI团队想用历史数据训练模型光做数据清洗就花了三个月。这些问题本质不是缺系统而是缺一个能统筹数据、让数据流转协同起来的“中枢”。所以“统好AI数据中枢”这个名字里的“统好”二字目标很清楚统一数据、好好协同。这里的“协同”不只是把数据搬到一起而是要让不同来源的数据在语义、血缘、时效、权限等多个维度上形成一致的、可被AI直接调用的协同状态。1. 项目最核心的设计思路为什么是数据中枢而不是中台行业内“数据中台”的概念已经谈了很多年但这个项目我们刻意用“数据中枢”来定义差异不在字面上而是在设计粒度上值得先说清楚。1.1 数据中心和中台的本质区别数据中台通常强调“共享”把数据资产沉淀为统一的服务能力由各业务方按需申请。它本质上是一个企业数据资产的“仓库市场”侧重供给端的统一。而数据中枢更强调“调度”它不仅是数据的存放场所更是数据流动、计算编排、AI调用和反馈回收的枢纽。我举个类比数据中台像一个大水库各条河流业务系统把水汇进去下游用户来取水。而数据中枢更像城市供水调度中心不仅管存量还要管每天哪个片区需要多少水、水质如何、压力多大、管道哪里堵了甚至还要根据天气预报提前调节水位。建设中台解决的是“有没有水”的问题建设中枢解决的是“水用得顺不顺、准不准、快不快”的问题。在实际项目中这个理念直接决定了架构设计。中台项目通常会先做数据分层ODS、DWD、DWS、ADS然后把重点放在数据模型和指标体系建设上。中枢项目在分层之上还要多做一层跨域数据协同层。这一层负责数据版本管理、血缘追踪、跨系统任务编排以及与AI模型的双向数据流转。这也是这个项目区别于普通数仓建设的最重要特征。1.2 “数据协同”到底在协同什么很多文章喜欢把数据协同挂在嘴边但落到实际场景我们梳理出来的是三个维度第一横向协同。指不同业务系统之间的数据打通。比如MES系统的生产批次数据、ERP系统的物料库存数据、QMS系统的质检数据在离散制造场景里需要按“工单号批次号”关联在一起才能完整还原一个产品的生产过程。没有协同这三套数据分属三个部门、三套口径出了质量客诉都很难追溯。第二纵向协同。指从原始数据到指标、再到AI模型输入特征的分层递进过程。原始数据是散乱的经过清洗、标准化、特征工程之后才能被统计模型或大模型使用。这个过程中任何一个环节的口径漂移都会导致最终AI结论失真。纵向协同就是把从ODS到特征库之间的转变过程管起来。第三人机协同。指业务人员通过自然语言或低代码方式向数据中枢发起数据请求中枢调用AI模型生成SQL、取数、生成分析报告。这里面有一个很容易被忽略的环节模型生成的结果需要反馈回中枢由血缘系统记录“这个结论用了哪些表、哪些指标”否则完全无法审计。这个项目最花时间的就是把这三个维度协同机制打通后面展开讲。2. 整体架构与关键设计选择分五层把AI纳入原生数据链路项目整体走的是“数据AI一体化”路线。我们没有把数据平台和AI平台分开建设而是把AI能力作为数据中枢的原生组件嵌入链路。整体架构分为五层每一层的定位和设计考量如下。2.1 五层架构从接入到服务第一层是接入层。负责对接各类数据源包含数据库MySQL、Oracle、PostgreSQL、消息队列Kafka、RocketMQ、SaaS API和文件存储。这一层我们统一封装了基于Debezium的CDC增量同步组件和基于SeaTunnel的批同步组件支持实时和离线双通道接入。第二层是存储与计算层。底层采用“湖仓一体”方案数据湖用Iceberg数据仓库用Doris和StarRocks。实时流处理使用Flink批处理使用Spark。这里的关键选择不只是性能而是Iceberg的表结构与Doris之间可以互相兼容查询省掉了大量传统“湖转仓”的ETL工作。第三层是治理层。包含元数据管理、数据质量、数据血缘和数据权限四块。其中元数据管理不仅管技术元数据库表字段还管业务元数据指标口径、业务定义。这个设计对于后续让AI理解数据非常关键我在2.3节单独讲。第四层是协同层。这是整个中枢的核心模块包含数据版本管理、跨域任务编排、调度依赖DAG、协同事件总线。它负责把前三层的数据能力组织成可协同、可编排的服务同时向第五层开放接口。第五层是智能服务层。包含三类能力一是指标问答基于自然语言生成SQL并查询结果二是知识库问答把企业文档和结构化数据做向量化后供大模型检索增强生成RAG使用三是预测与推荐面向营销、供应链场景提供标准化模型服务。这个层面我们用了Spring AI作为上层应用框架模型部署使用vLLM和Ray serving。2.2 选型背后的考虑为什么不全上顶尖组件很多技术团队做架构设计时容易陷入“什么最火用什么”的误区。这次我们在部分组件上故意选择了看起来“不够前沿”的方案原因只有一个团队可维护性。举个例子实时数仓很多人会选Paimon或Hudi但我们考虑到团队对Iceberg的运维经验更成熟、社区支持更完善最终仍然采用Iceberg。又比如向量数据库市面上选择很多我们最终使用了Milvus和PostgreSQL自带的pgvector混合方案。核心数据量在千万级以下用pgvector完全够用没必要再引一套独立服务减少一个组件就少一份运维负担。在“算力与存储”的规划上有一个参数计算值得参考。假设企业日增数据约1000万条单条数据平均2KB日增存储约20GB一个月约600GB。如果全部入湖保留冷热分层热存储保留90天约1.8TB冷存储持久化保存约7.2TB/年。按照目前市场上云对象存储约0.12元/GB/月的价格测算一年的冷存储成本约1万元左右。但如果把80%的数据都加载进Doris做实时分析那就是另一笔账Doris的本地存储成本约为云盘价格的数倍所以一定要把“访问频率低”的数据归档到Iceberg只把高频查询结果同步到Doris。这个冷热分离策略实际能为企业每年省下至少30%的存储成本。这块内容后文会结合代码详细拆解。2.3 统一建模与语义层让AI能“理解”数据这一层是这个项目里最具前瞻性的设计——语义层。我们引入了一个非常关键的思路给每张表、每个字段、每个指标都维护一套“语义描述”。比如“订单金额”这个指标技术口径是“订单表.实付金额”业务口径是“用户实际支付成功的金额含优惠券抵扣不含退款”。这套描述原本只有业务人员清楚AI模型无法直接读取。现在我们把技术元数据、业务口径、指标公式统一写入元数据中心配合向量化索引让大模型在生成SQL和回答问题时能检索到准确口径。具体实现上我们用OpenMetadata作为元数据底座扩展了自定义属性来记录业务口径。每张核心表在创建时强制要求填写数据字典未填写的表不会被AI问答功能引用。刚开始团队觉得这个流程太麻烦但当老板第一次用自然语言问出“上季度华东区退货率排名前五的产品是什么”而系统能准确给出包含业务口径解释的结果时所有人都觉得当初的强制要求是值得的。3. 数据协同引擎这是整个中枢的灵魂如果只做分层架构和元数据管理我们这个项目本质上还是数据平台。真正让“统好AI数据中枢”名副其实的是围绕协同而设计的核心模块。它包含四个方面其中几个环节是我们反复迭代了很多版才跑通的。3.1 数据版本与血缘协同的前提是可追溯先说一个场景。BI报表显示昨天销售额环比下降了12%老板问为什么。这时候数据分析师需要快速确认两件事一是数据本身有没有更新延迟二是口径有没有被改过。如果没有数据版本管理和血缘追踪排查这个问题可能要花半天。我们实现了一套基于Iceberg的Time Travel和自定义血缘采集器的方案。每次数据写入都自动生成一个数据版本快照血缘采集器通过解析Flink和Spark的执行计划自动生成表与表、字段与字段之间的依赖关系。这样任何一个指标变化可以通过血缘反查到上游原始表也可以正向推导出影响范围。这块的实际体验是在系统上线运行3个月后我们做过一次数据模型重构通过血缘分析直接锁定了受影响的15张下游表和应用系统提前通知相应团队回归测试整个过程没有出现一次数据质量问题投诉。做数据平台如果没有血缘管理相当于开车没有仪表盘速度再快心里也没底。3.2 跨域任务编排把不同团队的数据任务变成一条流水线在企业里数据任务通常是分布式的数据工程师跑清洗任务、算法工程师跑特征任务、BI工程师跑指标计算任务这些任务分属不同团队、使用不同调度器、依赖关系不透明。跨域协同要解决的就是把这些任务编排成一个大DAG统一管理依赖、优先级和告警。我们选择DolphinScheduler作为统一调度平台把所有数据任务从原系统的定时器迁移过来。核心配置是定义任务依赖表达式比如离线特征任务必须在ODS层同步完成后才能启动而指标计算任务又必须等特征任务完成。一个值得分享的参数调优经验是调度时间窗口的设置。假设ODS层日同步任务凌晨2点完成特征层任务需要40分钟跑完指标层任务需要20分钟跑完如果全部采用固定时钟触发会导致大量空等。我们改成事件触发模式上游任务成功通知下游任务启动整体链路耗时从原来的固定时间窗口的3小时压缩到1.5小时以内。这里的关键是建立协同事件总线任务节点通过MQ发布“任务完成事件”下游订阅事件后再触发而不是传统Cron表达式硬排。3.3 与AI Agent的协同让模型“会找数据、会取数据、会反馈”很多企业做大模型应用时会遇到一个尴尬模型很聪明但连不上企业数据。要么是让用户把数据导出来再上传给模型要么是给模型接一个查询接口但数据权限完全没法管控。我们这个项目的智能服务层通过AI Agent模式解决了这个问题。整体流程是这样的用户通过聊天窗口提问Agent接收请求后先调用意图识别模块把问题拆解为意图分类、实体抽取和参数提取三部分然后调用语义检索从元数据中心找到匹配的表和指标接着由Text-to-SQL模块基于检索到的表结构生成候选SQL最后执行SQL并返回结果同时对生成SQL的表增加权限校验用户只能查询自己权限范围内的数据。这里最容易出错的地方是Text-to-SQL的准确率。我们测试初期模型生成的SQL只有不到60%能直接跑通很大原因是表名和字段名的术语与大模型训练语料不一致。后来我们引入了“表结构业务描述历史问答对”三者联合构造Prompt的方法准确率提升到85%以上。另外还加了SQL语法预检在真正执行前先用解析器做语法校验避免生成的高风险SQL直接打到生产库这个前端校验设计非常关键。3.4 一次零售场景的完整协同走查光讲模块设计比较抽象分享一个完整场景某零售企业要做“门店智能补货”。原来的流程是总部运营导出各门店销售数据供应链团队用Excel做分仓预测再人工录入ERP生成采购单整个过程是割裂的全流程需要1到2周。接入数据中枢后链路变成了门店POS流水实时通过CDC接入KafkaFlink实时聚合各门店SKU销量数据落Iceberg的同时同步到Doris生成日维度的销售宽表。AI预测模型每天凌晨从协同引擎拿到特征数据输出补货建议表再由协同事件总线推送给ERP系统生成采购订单草稿。业务人员在系统中确认后即可完成补货。这中间的每一环都由数据中枢协同数据血缘完整记录从“POS流水”到“补货建议”的完整链路。最终效果是补货计划周期从10天缩短到1天库存周转天数下降了11%缺货率下降了近5个百分点。这就是“数据协同”四个字的具体价值。4. 核心落地过程从环境准备到链路贯通有了架构和协同引擎设计下一步就是实操落地。这一章我把能直接抄作业的内容写出来包括技术栈清单、核心配置和成本测算。4.1 技术栈清单与版本选择以下是这个项目实际使用的核心组件版本号以我们部署时为准模块组件版本用途实时同步Debezium Kafka2.3 / 3.5数据库CDC实时增量接入批量同步SeaTunnel2.3.3离线大表初始化与SaaS数据拉取流批计算Apache Flink / Spark1.17 / 3.4实时计算、离线批处理数据湖Apache Iceberg1.3.0存算分离的湖存储数据仓库Doris2.0加速查询与报表分析元数据OpenMetadata1.1技术元数据与业务口径管理调度DolphinScheduler3.1跨域任务统一编排调度向量库Milvus pgvector2.3 / 0.6AI检索增强的数据存储模型服务vLLM Ray Serve0.2.7 / 2.8大模型推理与在线服务AI应用框架Spring AI0.8封装AI Agent调用链路这里特别提一下SeaTunnel在离线初始化中的价值。老系统存量的历史数据通常有几十亿条如果全部走CDC同步会产生巨量 binlog 回放压力。我们的做法是先用SeaTunnel做全量初始化一次性把快照数据写入Iceberg然后开启Debezium的增量同步。全量增量的组合在实际执行时把数据接入风险大幅度降低对生产系统的压力也几乎可以忽略。4.2 数据接入链路的配置参考以MySQL通过Debezium同步到Kafka为例一个简化版的核心配置如下这个配置在生产环境经过验证适用于中等写入压力的业务库{ connector.class: io.debezium.connector.mysql.MySqlConnector, database.hostname: 10.0.1.10, database.port: 3306, database.user: cdc_user, database.password: xxxx, database.server.id: 54001, topic.prefix: ods_erp, database.include.list: erp, table.include.list: erp.orders,erp.inventory, snapshot.mode: when_needed, max.batch.size: 1024, poll.interval.ms: 300 }有一个参数容易被忽略database.server.id必须在整个MySQL主从复制环境中唯一否则会和其它同步任务冲突导致同步中断。这个坑我们在联调阶段踩过排查了大半天才发现是两个环境用了同一个server id。接着让Flink从Kafka读取变更数据写入Iceberg表核心作业配置如下CREATE CATALOG iceberg_hive WITH ( type iceberg, catalog-type hive, uri thrift://10.0.2.20:9083, clients 5, warehouse hdfs://nameservice/user/iceberg/warehouse ); CREATE TABLE iceberg_erp.ods.orders ( id BIGINT, order_no STRING, customer_id BIGINT, amount DECIMAL(12, 2), status STRING, updated_at TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) PARTITIONED BY (days(updated_at)) WITH ( format-version 2, write.metadata.previous-versions-max 50 );这段SQL里面我特意设置了write.metadata.previous-versions-max为50目的是在数据版本管理时能通过Iceberg的Time Travel回溯历史数据。但如果设得太大元数据文件会膨胀我们测试后认为50个版本在存储成本和回溯能力之间是比较合适的平衡点。4.3 算力与资源成本预估示例很多团队在项目汇报时会被问到一个问题这个平台要投多少钱给一套可复用的估算方式。假设企业日增数据量为1000万条单条数据平均大小为2KB那么日增原始数据约20GB。计算资源方面以Flink实时处理为例单并行度每秒能处理5000条变更事件1000万条日增分布在白天12小时峰值每秒约230条实际只需要1到2个并行度即可满足预留3倍冗余设置4个并行度足够。对应分配Flink TaskManager 4个Slot每个2核4GB总计算资源约8核16GB按常见云主机价格估算约每月2000元以内。离线批处理每天跑一次输入是全量历史数据假设累计1亿条Spark作业需要执行的聚合运算规模在10GB级别分配4个Executor每个4核8GB运行时长在20分钟以内。这部分成本较低。大模型服务是隐性成本大头。我们自建了一套7B参数的底座模型使用4张消费级GPU做推理和向量化推理吞吐约每秒30个token请求。如果要支撑全公司500人日常使用建议配置8张卡做负载均衡。如果预算有限可以做模型蒸馏用一个大模型生成训练数据蒸馏一个3B小模型到CPU上推理响应时间会慢一些但成本能降低约60%。这是AI infra方向比较省钱的打法。4.4 组织层面如何保障落地技术之外这个项目最容易被低估的是组织协同难度。数据中枢项目一定会触及各部门的利益边界业务部门担心数据被透明化后暴露出流程问题IT部门担心新增一套系统增加运维负担。我们的经验是从项目第一天就把各业务线的数据接口人拉进项目组每个业务线指定一名“数据Owner”由他们对数据口径和变更审批负责。中枢平台只提供技术实现不替业务方做口径决策。这样既守住了平台中立性也让业务方有参与感和责任感。上线后每次指标口径变更都通过协同事件总线通知到所有下游消费者避免业务口径变了但看报表的人不知道。5. 实战中遇到的典型问题与排查技巧这个项目从搭建到上线跑稳一共花了5个月中间踩过的坑不少。整理成一套速查表给正在做类似项目的同行省点时间。5.1 六大高频问题及排查方法问题现象可能原因排查思路解决方法CDC同步出现断点server id冲突或max_allowed_packet过小先查看Debezium日志中是否有重复server id报错为每个连接器分配唯一server id调大MySQL的max_allowed_packetIceberg小文件过多Flink checkpoint间隔过短、写入频率过高检查表目录下文件数和平均文件大小调大checkpoint间隔到60秒开启compaction报表数据与其他部门口径不一致缺少业务口径统一管理对比两边的指标公式和过滤条件对指标建立口径字典强制由语义层统一输出AI问答返回空结果Text-to-SQL生成的表不在用户权限内排查权限表配置和血缘元数据是否完整对核心表统一打权限标签设置默认白名单调度任务互相阻塞共享资源池被大任务占满查看DolphinScheduler队列和运行中的任务为大任务单独划分资源队列设置任务优先级模型服务响应变慢并发请求突增导致GPU排队查看模型服务监控的请求队列长度增加推理副本启用请求排队降级策略5.2 排查链路问题的三个实用命令日常运维中需要快速定位问题根源下面三个命令是我们在项目中最常用的。第一个是查看Kafka消费组积压判断实时链路是否卡住kafka-consumer-groups.sh --bootstrap-server kafka01:9092 \ --describe --group flink_sync_group如果LAG列持续增长说明Flink作业消费速度跟不上优先检查CDC源端数据库压力再看Flink作业反压情况。第二个是直接查询Doris中的表分区命中情况确认查询是否走到了正确的分区SHOW PARTITIONS FROM ods_erp.orders;如果查询扫描的分区数量过多说明分区字段或过滤条件设置不合理这时候需要回到表中检查是否用日期字段做分区。第三个是从Iceberg快照查看历史版本用于恢复误操作SELECT * FROM iceberg_erp.ods.orders. FOR SYSTEM_TIME AS OF 2024-06-01 00:00:00;这条SQL能直接读取历史版本数据。我们在一次误删全表的备份恢复时就是靠它救回来的所以建议所有入湖的表都开启Time Travel保留至少30个版本。5.3 三个减少返工的设计习惯第一建表时强制添加updated_at字段。几乎所有数据同步、变化捕捉和审计需求都依赖这个字段没有它就很难做增量更新。即使业务系统表本来没有也必须在入湖时补上。第二所有同步任务必须加上数据质量校验算子。我们写了一个通用的数据量波动检测组件如果当天同步的数据量与近7天均值偏差超过20%直接阻断下游任务并告警。这个设计在几次源头系统改造导致数据丢失的事件中起到了关键拦截作用。第三从第一天就做权限模型设计不要等上线再补。数据中枢越往后接入的角色和系统越多权限收敛的难度成倍增加。我们采用的是基于标签的权限控制给每个用户、每张表都打标签通过标签匹配来控制可见范围。改造起来比写死用户表灵活得多。6. 数据协同带来的可见变化最后用一个项目上线后的真实对照看看数据协同到底给企业带来了什么。指标项目前项目后核心经营数据出报表周期T2天T0实时跨部门数据口径一致率约65%98%以上数据分析需求平均交付周期7天1天内AI问答回答准确率未启用85%以上数据任务运维告警处理时长3小时30分钟这些数字背后不只是单纯的技术升级更是工作方式的改变。以前业务人员要数据靠“求人”在群里喊半天等数据工程师排期。现在直接在对话窗口问产品数据中枢通过AI Agent自助取数合规、快速、口径统一。数据团队终于能把精力从取数需求中解放出来去做更深层次的数据分析与模型优化。在实际项目复盘时我最大的体会是做数据中枢这类系统最难的不是技术选型或架构设计而是让不同团队真的愿意把数据交出来、把口径对齐。技术平台能解决“能不能协同”的问题但“愿不愿意协同”取决于项目治理机制和组织推动力。如果有同行正在做类似项目建议一开始就花1/3的精力在管理机制和业务方关系上否则再优雅的架构也很难真正跑出价值。最后分享一个小技巧协同层的事件总线是整套系统最容易被忽略但又最重要的组件。很多团队把大量精力花在存储和模型上但真正让数据流动起来、让任务串成链路的是那条看似简单的事件通道。把事件总线设计好、监控好数据中枢就成功了一半。
返回列表