
从云栖大会回来这几天我一直在琢磨一个信号AI应用好像终于不“卷”模型本身了。会场里大家嘴上聊的还是大模型、Agent但实际上蹲在展位前问得最多的是实时数据怎么接进来、脏数据怎么处理、链路延迟能不能压到秒级。这正好印证了今年云栖的一个核心趋势——AI应用进入生产阶段大家开始拼「实时数据智能」。说白了模型能力已经成了门槛而不是壁垒真正决定一个AI应用能不能在线上稳定跑出价值、能不能扛住真实业务压力的是它背后的实时数据链路是否够快、够准、够干净。这篇文章想结合我在云栖现场的观察和一些实际项目经验把“实时数据智能”这件事掰开揉碎讲清楚给正在做AI应用开发、或者在运维侧跟数据链路打交道的朋友一些参考。1. 从“跑通Demo”到“撑住生产”AI应用的真实分水岭1.1 百人排队抢GPU的现场藏着行业风向云栖大会的AI展区今年有个很有意思的现象专门讲模型微调、Prompt工程的展台人气反而没有几个数据平台的展台高。我在一个实时计算相关的展位前站了十几分钟听到的几乎全是这类问题——“我们想做一个智能客服大模型已经接好了但用户的历史订单数据在MySQL里还有一部分行为日志在日志系统里怎么才能把这些数据实时喂给模型” “智能推荐上线之后用户刚点了一个商品下一次刷新推荐位还是老内容结果点击率掉了一大截你帮我看看是不是数据链路延迟的问题。”这些问题的共性非常明显模型能力已经不再是瓶颈大家默认大模型能搞定语义理解和生成真正卡住AI应用落地脚步的是数据。更准确地说是实时数据。这里的“实时”不是指毫秒级的流式计算而是要覆盖一个AI应用从拿到原始数据、完成加工清洗、再到被模型消费的完整周期。在展位上我没有直接回答他们但脑子里已经把这几年做数据平台和AI应用的经验串了一遍AI应用生产化本质上是数据基础设施的竞赛不是模型竞赛。1.2 模型能力有上限数据智能才是生产瓶颈做一个AI应用Demo阶段通常只要一个模型接口、一份测试数据就能跑通。但进入生产你要面对的是三种现实数据一直在变而且变得很快。用户的行为、库存状态、价格、内容热度都是分钟级甚至秒级变化。如果模型只能读到昨天甚至一小时前的数据那它给出的答案天然就是滞后的。数据质量参差不齐。线上数据有重复、有缺失、有时序错乱。脏数据直接会让模型输出错误结果而这个错误结果在Demo阶段根本看不到。数据规模远超单机内存的承受范围。你得考虑数据怎么存储、怎么分片、怎么保障查询性能这部分工作跟模型完全无关但又是生产环境逃不掉的事。我第一次帮客户把智能问答从Demo推到生产时印象最深的一句话是客户CTO说的“我们不缺模型缺的是能让模型用上最新数据的那条管子。”当时我们把日志系统里的用户行为数据通过实时链路接入到向量库再把向量检索结果拼进Prompt那套系统上线后效果提升非常明显。从那以后我就确信AI应用的生产力是由实时数据智能决定的。模型参数是你买得到的实时数据链路是你自己一点一点建出来的这才是真正的护城河。1.3 “Demo快生产难”的六个典型卡点如果非要把AI应用生产化的难点列清楚我觉得这六个卡点最典型。它们都是我在真实项目里遇到过的也是这次云栖现场被问得最多的问题卡点表现根因数据新鲜度不够用户刚发生的行为模型感知不到离线批处理链路 T1缺少实时同步特征与上下文缺失模型问了上下文系统答不上来业务数据分散在多个系统没有统一服务层查询性能扛不住检索或特征查询超时模型被迫降级直接查业务库OLTP 和 AI 场景混部数据质量不稳定模型偶尔输出离谱答案源头脏数据、重复数据没有清洗规则链路可观测性差线上出问题只能靠猜缺全链路 Trace 和数据血缘成本失控实时任务越来越多集群费用暴涨缺少资源治理和分层存储策略这六件事没有一件是模型层能解决的。它们全都依赖一套扎实的数据基础设施。换句话说AI应用进入生产阶段后研发团队真正的工作重心从“调模型”转移到了“调数据链路”这恰好也是我在云栖上感受到的最明显风向。2. 把数据变成AI能吃的“实时粮草”三个关键动作2.1 数据同步从“T1批量搬运”到“秒级流动”传统数仓时代数据处理的基本单位是“天”。每天凌晨跑一批任务把昨天的数据算成报表。这种模式对AI应用来说完全不够用。一个推荐系统如果只能用昨天的数据那用户昨晚刚加购的商品今天推荐位上是看不到的一个风控系统如果只能处理昨天的数据那今天发生的盗刷行为明天才能被拦截。生产级AI应用要求数据从产生到可被消费的延迟控制在秒级到分钟级这意味着原来T1的批处理任务必须替换为实时同步链路。当前最主流的做法是两类配合一类是用CDC工具比如Flink CDC把业务库的变更日志实时采集出来再写入消息队列或数据湖另一类是把埋点、日志、物联网设备上报等半结构化数据通过实时计算引擎进行流式清洗。我自己的建议是先别急着追求“全链路实时”而是优先保障对AI效果影响最大的那部分数据——用户实时行为、订单状态、库存变化——先打通这条实时管道其他低频数据可以继续走批量。一次接入太多实时源运维复杂度会指数级上升这跟做缓存的原则差不多先给最热的数据加缓存而不是给所有数据都加缓存。2.2 数据加工让原始数据变成模型能读懂的“上下文”实时同步只是把数据搬了过来原样的数据并不能直接被AI消费。举个例子用户行为日志里有的是“点击”事件有的是“曝光”事件还有“加购”、“支付”等等如果直接把这种原始日志喂给模型大模型做推断效果会很差因为混杂的事件流实在太噪。数据加工要做的就是把原始日志整理成结构化的特征序列。常见的加工动作包括清洗去重、过滤异常值、纠正错乱的时间戳。聚合把原始事件聚合成用户维度的行为序列比如“最近1小时浏览的商品列表”。特征计算统计类特征点击率、停留时长、序列类特征行为轨迹、实时标签当前情绪、意图。这些特征本质上就是“实时数据智能”的产物。加工这条链路最考验工程功底的是状态管理。流式计算里每个用户的会话状态、滑动窗口状态、去重状态都需要在内存里维护状态如何持久化、如何容错、如何保证精确一次语义直接决定了加工出来的数据准不准。Flink 在这块做得很成熟它的状态后端支持 RocksDB可以应对很大的状态规模。我在项目里用 Flink SQL 加工用户行为序列一个比较简单但实用的实时特征流水线大概长这样-- 示例计算用户最近5分钟的商品点击序列用于推荐场景的实时上下文 CREATE TABLE user_click_log ( user_id BIGINT, product_id BIGINT, action STRING, click_time TIMESTAMP(3), WATERMARK FOR click_time AS click_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic user_click_log, properties.bootstrap.servers kafka:9092, format json ); CREATE TABLE user_recent_click_seq ( user_id BIGINT, product_seqs STRING, window_start TIMESTAMP(3), window_end TIMESTAMP(3) ) WITH ( connector kafka, topic user_recent_click_seq, format json ); INSERT INTO user_recent_click_seq SELECT user_id, LISTAGG(product_id, ,) WITHIN GROUP (ORDER BY click_time) AS product_seqs, TUMBLE_START(click_time, INTERVAL 5 MINUTE) AS window_start, TUMBLE_END(click_time, INTERVAL 5 MINUTE) AS window_end FROM user_click_log GROUP BY user_id, TUMBLE(click_time, INTERVAL 5 MINUTE);这段 SQL 最核心的是把 Kafka 里的原始点击流转成了以“用户 5分钟窗口”为粒度的商品序列字符串模型接过去之后可以直接作为 Prompt 里的上下文。像这样的加工逻辑规模大了之后会非常多所以我在项目里一直强调能流上做的不要落到离线能SQL表达的就别写复杂代码因为后期要维护的逻辑越简单越可靠。2.3 数据服务把“数据”和“特征”变成API加工好的数据最后要能被线上AI服务调用。我们不可能让推理服务直接去读消息队列或者查数据湖那样延迟不可控而且业务逻辑和数据物理存储会耦合在一起。生产环境里正确的姿势是加一层“数据服务层”把加工结果封装成高并发、低延迟的API。这层服务做得好的话AI应用的数据获取就可以变得非常纯粹像这样推理服务调用“用户最近行为API”入参user_id返回最近1小时行为序列。推理服务调用“实时商品特征API”入参product_id返回商品当前库存、价格、热度。推理服务调用“相似人群API”入参user_id返回种子人群包ID列表。我在云栖的 DataWorks 展区看到现在不少平台已经在把“数据开发 数据服务”整合成一套完整链路数据开发完可以直接发布API甚至还能生成API文档。这对于很多中小团队来说挺友好的因为过去这层服务通常要自己写一套后端现在平台已经帮我们解决了很大一部分工程问题。但我也要提醒一句平台能帮你生成API却帮不了你设计API的语义和粒度。这个API是给AI服务用的不是给人写的管理后台用的所以响应要快、结果要结构化、更重要的是要带新鲜时间戳让模型层知道这份数据的生产时间。3. 一条生产级实时数据链路的搭建笔记架构、延迟与成本3.1 实时链路的标准组成以及“为什么是这套组合”这里我画不了一张完整架构图可以直接用文字描述标准链路的分段逻辑。一套生产级AI应用的实时数据链路大致会分成四段数据接入层负责从业务系统、埋点日志、外部数据源采集数据。采集技术通常是CDC或消息队列Kafka / Pulsar / RocketMQ。实时计算层负责清洗、关联、聚合、特征化。这层主要用 Flink 或 Spark Structured Streaming。存储与索引层负责为AI服务提供低延迟的数据访问。常见选型有实时数仓或OLAPStarRocks、ClickHouse、向量数据库Milvus、KV存储RedisHBase。数据服务层负责把存储能力封装成接口供推理服务调用。这套组合不是谁拍脑袋定的是经过生产环境反复锤打之后沉淀下来的“标准答案”消息队列解决的是削峰填谷的问题。实时数据往往是突发性的比如大促流量如果直接打到数据库很容易被打挂。队列能缓冲压力让下游按自己的速度消费。Flink解决的是状态化计算的问题。AI应用需要的数据加工大多涉及时间窗口和累计状态Flink的容错和状态管理在这一块做得很成熟。OLAP/向量库解决的是低延迟查询的问题。AI推理通常需要宽表式的多维特征OLAP的列存和向量化执行很适合这种场景。向量库则承接非结构化语义检索的场景比如知识库问答。服务层解决的是解耦和SLA保障的问题。不把数据存储直接暴露给模型方便做缓存、限流、熔断和权限管控。3.2 端到端延迟是怎么抠出来的我见过不少团队嘴上说要做到“实时”实际链路搭完跑起来端到端延迟却在分钟级。为什么因为每个环节都会悄悄增加延迟而且越往下游延迟被放得越大。我一般会把链路拆成几段去定位采集端延迟业务库CDC的binlog监听是从数据库提交事务后开始的通常毫秒级。但埋点日志的采集涉及客户端上报和网络传输可能需要几秒到几十秒。消息队列排队延迟Kafka在高吞吐下如果分区数设置不合理消费端Lag会飙高数据在队列里积压。计算延迟Flink的窗口和状态计算如果不是纯增量而是有join操作性能会下降。尤其是流流join比如点击流和曝光流做关联状态规模一大吞吐就掉。存储写入延迟OLAP适合批量写入频繁小批量写入会导致写入吞吐瓶颈。很多团队没做攒批优化实时链路变成了“慢链路”。有一次我们做智能客服的实时知识同步客户反馈“新上架的商品客服机器人5分钟后才知道”。排查下来卡点竟然是OLAP写入的攒批策略设置了60秒flush加上上游CDC的采集渠道本身有2秒延迟再叠加消息队列的批量拉取整体延迟就积累到了分钟级。调优方案很简单把攒批时间从60秒调成5秒同时把消费端的拉取批量参数调整了一下端到端延迟直接降到8秒左右。这类问题不亲自动手排查真的很难意识到延迟不是某一个点造成的是每个点你让了一寸最后合起来就是一大截。3.3 成本账实时链路很贵怎么把钱花在刀刃上实时链路比离线链路贵这是事实。原因很直接实时任务需要长期占用计算资源不能像离线任务一样跑完就释放状态存储、消息队列、OLAP 集群也都是持续运行的。我建议团队在做实时数据智能时先算一笔账环节主要成本优化手段采集数据库binlog读取对源库性能有损耗尽量用低峰期同步减少对在线业务影响消息队列Broker 节点与存储占用按峰值流量预留同时设置合理的留存时间实时计算Flink 资源常驻用 Flink 的自动扩缩容削峰填谷存储OLAP 节点多副本区分热数据/温冷数据冷数据落到对象存储服务层API 实例数配置合理的缓存策略避免重复计算这笔账算完之后第二个动作是分层把实时链路里真正关键的“黄金数据”比如核心交易、用户实时行为保留在毫秒级成本最高的链路上把次重要的数据放到秒级链路把不那么重要的数据每天批量同步一次完全够了。很多团队一上来就想“全实时”成本直接爆炸效果却没有翻倍。实时数据智能不是所有数据都实时而是让能产生智能化价值的那部分数据实时。4. 实时数据智能落地中的五个典型坑与排查思路4.1 “实时跑批”不是实时关于计算语义的大坑我第一次接触实时数据处理时犯过一个特别经典的错误把离线批处理的SQL直接原样搬到Flink上跑。结果看起来都在“实时计算”但数据算出来的结果总是对不上。原因很简单批处理是“读全部数据然后算一次”流处理是“数据一条条到来边来边算”两套语义完全不同。举个具体例子。统计“今日UV”离线批处理会对全天数据做去重结果精确实时Flink如果直接Count Distinct内存会被撑爆。而如果改用HyperLogLog近似去重结果看起来对但和离线对不齐。这就导致业务方拿着实时报表去跟离线报表做校验永远差一点最后怀疑实时链路不稳。这类问题的根子不在代码在于“实时计算的结果到底允许多大的误差”这件事没有跟业务方达成共识。后来我们在项目里定了规矩凡是涉及精确去重、精确累计的指标尽量用窗口状态做精确计算凡是能接受近似值的指标比如用户画像的活跃度分桶才用近似算法。业务口径对不齐链路再快都是白快的。4.2 写入高频与查询高频打架OLAP选型要分层很多AI应用场景里实时特征服务既要高并发查询又要高频写入。比如需要实时写入用户行为指标同时又要供在线推荐服务高并发查询。如果把这两类负载放在同一个ClickHouse节点上非常容易互相干扰写入一多查询P95就飙升查询重了写入又会阻塞。这不是ClickHouse不行而是架构上写查询混合的问题。我的经验是分两层第一层用类StarRocks或ClickHouse承接单条/小批量的写入和近实时查询做数据分析和后端的查询计算第二层用Redis等KV存储承接高频在线查询把计算结果比如“当前用户最近浏览的N件商品”定期或者实时推入缓存推理服务只查缓存不查OLAP。这个分层方式的核心逻辑是OLAP负责算Redis负责扛。算是一种能力扛是一种吞吐两者混在一起多半会出问题。如果你做的是知识库类AI应用向量检索可能也需要单独一套Milvus更不能和大宽表的OLAP混用否则那一侧的查询性能也会出状况。4.3 数据血缘缺失AI应用排障变成考古现场AI应用链路比传统应用长得多数据源、消息队列、Flink任务、OLAP表、API服务、模型推理六七层链路。只要其中任何一层的数据语义变了比如前端埋点把商品ID的字段类型从String改成了Long下游所有层都会静默出错。最难受的是模型不会报错——它只是开始给出更离谱的答案。我们曾经排查过一个智能问答系统“突然变笨”的问题查了很久才发现原来上游埋点改了字段名导致Flink解析出来的商品ID全是默认值知识库召回结果全乱了。因为没有血缘关系图排障过程简直是在考古翻Flink任务、翻Kafka Topic、翻埋点文档、翻数据仓库表结构最后才找到源头。从那以后我们在所有链路都强制维护血缘信息。现在一些平台比如DataWorks已经内置了比较完善的血缘功能能自动解析SQL任务的上游下游依赖查问题的时候顺着血缘图一路点下去定位会快很多。对于还没上平台的自建链路至少要在任务命名、数据表注释、API文档里做足审计标记别等出问题再去翻。4.4 回放与补偿实时链路最容易被忽略的一环实时链路的特点是“无法回头”。离线报表算错了可以重跑实时任务出了问题错过的数据就是错过了。比如Flink任务因为上游Kafka数据格式异常从10:00到10:05连续五分钟反序列化失败这五分钟的数据如果不做补偿就永久缺口了。这在AI应用里很致命——少了一段高峰期的行为数据模型推荐出来的内容就可能完全走偏。处理办法我一般会给实时链路加“旁路”原始数据在进入计算引擎之前按原始格式也备份一份。一旦发现问题可以写一个临时的批处理任务把缺失时间段的数据从备份里重新刷到下游。这种设计成本很低但救过大命。在做实时数据智能时我总跟团队强调实时链路的设计目标不只是“快”更要是“可恢复的”。否则再快出一次事故你前面赚到的效率都会补回去。4.5 别让“实时”变成业务方的焦虑最后一个坑不是技术问题是预期管理问题。实时能力一上线业务方就会天然期待所有数据都是实时的。但业务和技术的现实是有些数据天生就是“非实时”的比如财务稽核指标、需要人工审批的状态字段。这个时候你硬要做到实时只会花了大成本业务方还不满意。我的做法是对外明确“数据新鲜度分级”。比如交易风控场景要求秒级用户画像要求分钟级报表统计可以小时级。然后针对每一级数据定义好SLA写清楚数据延迟上限。别让业务方拿着“实时”两个字当成万能尺子量一切你自己的链路也能少一点无效的优化。这个看似不是技术问题反而是我从多个项目里总结出来最值钱的经验。5. 云栖上的风向以及每个工程师该补的课5.1 DataWorks的方向数据平台向“AI原生”迈进这次在云栖DataWorks给我最直观的感觉是它已经不再只是一个“数据开发IDE”而是开始贴着AI应用开发流程做了很多原生能力。比如在数据开发过程中直接集成了智能补数据、智能建模的能力甚至能将加工好的数据发布成API给AI引用走通“数据集成-数据开发-数据服务-AI应用”的完整闭环。对开发者来说这意味着过去需要自己拼装的“数据管道全家桶”现在被整合在了一个平台里。当然要说明一下这类平台能力更适合以数据中台建设为主要诉求的团队如果你的场景特别垂直、数据量特别大可能还是需要定制化自研。但平台解决的“血缘追踪、统一权限、监控告警、运维管理”这些问题自研的成本真的非常高中小团队完全没必要重复造轮子。5.2 “AI Agent”开始吃实时数据从玩具走向生产力本届云栖上“智能体应用”的提法非常多比如扣子Coze这类平台让普通人也能搭出相对复杂的Agent。但我们做工程的需要留意一个区分搭一个Demo Agent很容易做一个能稳定完成业务流程的Agent非常难。难在哪儿难在Agent需要“感知环境”。它得知道订单当前状态是什么库存还有多少用户刚才做了什么甚至用户当前情绪如何。这些信息没有一样能从模型参数里“猜”出来全都需要实时数据供给。我判断接下来半年到一年会出现一批专为Agent服务的“记忆层”和“情境层”基础设施。它们的作用是把实时业务数据组织成Agent能检索、能消费、能感知的格式这其实就是“实时数据智能”在Agent场景的落地。如果现在你正在开发Agent应用我的建议是可以先把注意力从研究“怎么设计更好的Prompt”移到“怎么让Agent拿到此时此刻的数据”上——后者带来的效果提升通常比前者大得多。5.3 运维工程师和AI应用开发者的角色变化之前有人问我说AI大模型应用开发是不是只需要懂Prompt不需要再懂数据工程和运维了我的观点是门槛确实降低了但竞争力门槛反而提高了。以前只写SQL、只调参数的那波人如果不懂数据从哪里来、链路怎么监控、故障怎么恢复在AI应用生产化的阶段会很难受。反过来今天最吃香的工程师往往是既懂模型调用逻辑又懂实时数据管道还能处理数据质量问题的“T型人才”。对运维工程师来说AI应用上线之后监控的不再只是CPU、内存、QPS而是要监控“数据新鲜度”“链路延迟”“丢失率”这些数据语义层面的指标。这要求运维工具和知识结构都得更新。在云栖现场我也看到和AI应用可观测性相关的讨论明显比去年多了。趋势已经很清楚AI应用开发会越来越往工程化走而工程化的核心正是数据智能。跑了几届云栖今年最大的体感是AI应用能不能成模型好不好只是下限数据链路的实时性和质量才是上限。模型会越来越开放、越来越便宜最后真正拉开差距的是谁能稳定、实时、高质量地把数据喂给AI。如果你也在做AI应用我真诚建议分一部分精力去打磨数据管道而不是只盯着模型本身。这条路我是从一次次生产事故里趟出来的希望你少踩几个坑。