拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Doris + Paimon 构建 Agentic AI 数据闭环

最近在帮团队搭一个 Agentic AI 项目的数据底座老板上来第一句就是“把大模型 API 接上、工具链配上是不是就完事了”我直接打断他还差最重要的一环——数据闭环。Agent 不是简单的“模型调用”它每次感知、决策、调用工具、产生反馈都会留下轨迹。这些轨迹如果只是散落在日志里Agent 就是一次性问答机器只有被统一采集、存储、分析并回流到下一次决策它才会越用越准。我们最终落地选择了 Apache Doris Apache Paimon 2.0 这套组合Paimon 2.0 负责把全量数据“存厚”Doris 负责把关键数据“读薄并加速”。这篇内容就把我们的整套架构、实操步骤和踩过的坑一起讲清楚。1. 先想清楚Agentic AI 需要什么样的数据闭环1.1 Agentic AI 的核心不只是模型Agentic AI 和传统程序最大的区别是它具备“感知—决策—行动—反思”的循环能力。一个采购 Agent 要查历史报价、对比供应商、生成订单、跟踪物流每一步都会产生结构化动作和结果一个人力 Agent 要处理简历筛选、面试邀约、候选人反馈同样会留下大量中间状态。这些中间状态就是数据资产。问题是大多数项目把 Agent 的会话记录和工具调用日志丢给 Elasticsearch 或者普通消息队列用完就扔。短期看没问题一旦要做效果评估、故障回溯、prompt 迭代、模型微调就会发现根本没有可复用的数据层。我见过不少团队Agent 上线后只能靠肉眼抽查聊天记录来判断好坏这显然走不远。所以 Agentic AI 对数据平台的要求和传统 BI 完全不同。传统数仓更关心“订单量涨了多少”“转化率怎么样”Agentic AI 更关心“某一次决策的完整轨迹是什么”“如果换一个动作结果会不会更好”。这要求数据能够被回放、被更新、被关联并且要在一个相对短的周期内流回决策链路。1.2 数据闭环的四个关键环节我习惯把闭环拆成四个环节采集、沉淀、反馈、进化。采集是基础。Agent 每一步的输入输出、工具调用参数、返回结果、环境状态、用户反馈都要记录成事件。注意这里不能只记回答内容还要记上下文、token 消耗、延迟、RAG 命中情况、工具调用顺序否则后面做成本分析和效果归因会非常痛苦。沉淀是关键。事件数据进到存储层之后不能只是一堆 append-only 日志。很多数据是需要更新的比如一条轨迹的奖励分可能在整轮任务结束后才确认Agent 的记忆状态也会随新交互而变更。这就要存储层支持“流式更新批量分析”并存。这也是我们选择 Paimon 2.0 的核心理由。反馈是灵魂。没有反馈的闭环不叫闭环。用户点没点赞、任务完没完成、业务转化有没有发生这些信号必须回流到数据系统。反馈可以是人工标注也可以是业务系统自动推送给结果。进化是终点。离线分析这些反馈生成新的特征、修正 prompt、更新 RAG 知识库、微调模型然后把结果重新部署到 Agent 运行环境完成整个飞轮。只有把四个环节串起来Agent 才能持续变聪明。1.3 怎么判断闭环是否合格判断一个闭环是否合格不要只看“有没有循环”要看四个指标数据新鲜度从 Agent 产生事件到数据可被决定使用需要多少秒轨迹覆盖率是不是所有 Agent 动作都有完整 trace能不能任意回放反馈回填率真实业务结果有多大比例回流到了数据平台样本产出周期从“新反馈出现”到“训练样本可用”需要多久我见过很多系统离线训练用 A 流程在线特征查询用 B 流程两者口径还对不上。闭环最终要的是“一个事实源头多个使用入口”。Paimon 负责事实源头Doris 负责多个使用入口中的低延迟查询。下面详细说为什么这么组合。2. 为什么选 Apache Doris Paimon 2.0不是巧合2.1 Paimon 2.0给数据闭环一个能流式更新的湖Apache Paimon 是流式数据湖存储底层以列存格式保存数据天然适合和 Flink 搭配做流式入湖。2.0 版本在流更新、主键表性能、增量读取、自动小文件治理等方面都有明显改进。Paimon 2.0 最强的一点是它既能存 append-only 事件流也能存需要主键更新的状态表并且下游可以通过 Flink 读取到变更日志。这个能力对 Agentic AI 太重要了。还是说采购 Agent初始决策可能是“选价格最低的供应商”但后来线上反馈“这家供应商发货经常延迟”我们需要更新这条决策的奖励分和供应商评分。如果只靠离线批量覆盖在线 Agent 永远在用旧数据。Paimon 2.0 的主键表可以让奖励分、评分这类数据在流式链路里被实时更新同时保留历史快照方便事后回溯。另外Paimon 2.0 支持增量读取。这意味着我们不用每次全量扫描湖表可以直接读“从某个 snapshot 之后变化的数据”用来触发实时特征更新或告警。这一点对跑 Agent 反馈回流很关键我在 4.3 里会再说。2.2 Doris给在线智能一个低延迟查询出口Apache Doris 是分布式 MPP 分析型数据库它能直接通过 Multi-Catalog 查询 Paimon 中的数据也可以把热数据导入本地表做高性能分析。Agent 运行时会调用大量特征查询例如“这个用户最近 3 天和这个供应商交互过几次”“这个动作的历史成功率和平均耗时是多少”。这种查询的特点是并发不低、延迟要求高、过滤条件多。Doris 有物化视图、前缀索引、布隆过滤器、分区裁剪等能力配合向量化执行引擎能在亚秒甚至毫秒级返回聚合结果。我们实测下来大部分 Agent 在线决策查询的 p95 能控制在 50ms 以内前提是数据模型和 SQL 写得合理。另外Doris 还支持高并发点查。虽然很多团队会把“Agent 的记忆”放到向量数据库但向量库往往只适合做相似度检索不适合做条件过滤和统计聚合。Doris 可以对向量元数据做过滤比如“先筛出最近 7 天的有效记忆再做 embedding 匹配”这种组合查询交给 Doris 更合适。2.3 组合的分工边界Paimon 和 Doris 的分工用仓库和分拣台类比更清楚。Paimon 就是把所有货品全量轨迹、反馈、状态变化原样存在仓库里支持随时盘点和按版本回溯Doris 是分拣台把最常用的货放到手边快速响应决策需求。但分拣台不是仓库不可能把全部历史都堆在上面所以要设计冷热分层。我们实际的边界是全量数据进 Paimon 2.0作为事实源头保留时间窗口尽量长。在线决策需要的特征和统计结果物化到 Doris 本地表或物化视图。热数据和冷数据之间通过定时任务或流式任务同步而不是把 Doris 当唯一存储。所有离线训练样本、大规模数据挖掘直接读 Paimon 或把它导出到训练平台避免干扰在线查询。这个组合最大的价值是“一套数据两个入口”。离线分析和在线查询不再各搞一套闭环才能稳定跑起来。3. 一套可落地的架构设计3.1 闭环的两个循环在线小回路与离线大回路整个架构里其实有两个循环在跑。第一个是“在线小回路”Agent 触发决策时从 Doris 查询特征Agent 执行动作后产生事件事件经过消息队列进入 PaimonPaimon 的增量数据再回流到 Doris更新后续特征。这个回路的目标是分钟级甚至秒级闭环适合那些需要快速响应的信用评分、风险拦截、动态推荐等场景。第二个是“离线大回路”Paimon 里积累一天或一周的数据经过离线分析生成训练样本、特征字典、知识库内容模型或 prompt 更新后发布到 Agent 运行环境Agent 新一轮表现产生的反馈再进入 Paimon供下一轮优化。这个回路通常以小时或天为单位适合系统性提升 Agent 能力。两个循环不是互相独立的。在线小回路负责“执行效率”离线大回路负责“系统进化”。共用 Paimon 作为数据底座就不会出现在线日志和离线仓库对不上号的问题。3.2 数据模型怎么设计数据模型设计直接影响闭环能不能跑通。我们在 Paimon 里主要建了四类表轨迹事件表记录 Agent 的每一步动作和观测结果。分区按天字段包括 session_id、event_id、agent_id、user_id、action_type、action_content、observation、timestamp。这类表以 append 为主不设主键写入性能优先。记忆状态表记录 Agent 对某个实体或关键词的持久化记忆。主键表主键可以是 agent_id entity_key key_name字段包括 value、last_update_time、expire_time。这个表会被在线 Agent 高频读取。反馈评分表记录人工或业务系统返回的反馈结果。主键表主键是 session_id step_id字段包括 user_score、business_result、refund_flag、comment。特征聚合表Doris 侧的物化结果比如“某供应商近 7 天平均评分”“某动作成功率”用于在线决策查询。分区和主键设计有几个容易踩的坑。Paimon 主键表如果带分区分区字段必须能由主键推导出来否则跨分区更新会失败。比如把 event_time 作为分区字段但主键只有 event_id那主键表就没法确定事件属于哪个分区。我们的做法是把 event_time 也放进主键或者用“业务日期”做分区并让主键包含它。3.3 闭环真正“合上”的地方用一个具体场景描述闭环用户让 Agent 推荐一家本地保洁公司。Agent 先通过 Doris 查询候选公司在过去 30 天的订单量、好评率、平均响应时间然后调外部服务获取实时价格最后综合这些信息给出推荐。用户选择了其中一家或者给了差评。这些行为进入 KafkaFlink 消费后写入 Paimon 的轨迹表离线任务扫描新的反馈更新每家公司的评分更新后的评分写回 Paimon 或 Doris 的特征表。第二天另一个用户再问相同的问题时Agent 会拿到新的评分数据从而改变推荐结果。这就是闭环真正合上的位置不是某一条数据管道跑通了而是“分析结果能反过来影响下一次决策”。Doris 和 Paimon 分别承载了这个环的两端。4. 实操打通 Doris 与 Paimon 的关键步骤4.1 环境和版本准备先说版本。我们用的是 Doris 2.1.4Paimon 2.0.1Flink 1.18对象存储用的 MinIO 模拟 S3。Doris 从 2.x 开始对 Paimon 的 Multi-Catalog 支持比较完善。如果你用老版本 Doris建议先升级否则可能在读取 Paimon 2.0 的表格式时遇到兼容性问题。Paimon 本身不依赖 Hadoop 集群只要有文件系统或对象存储就能跑。生产环境可以用 S3、OSS 或 HDFS。我们在本地测试时直接用 MinIO配置起来最省事。注意 Paimon 的元数据默认放在 warehouse 目录下所以 Doris 挂载时只需要指向 warehouse 根路径。4.2 在 Doris 中挂载 Paimon CatalogDoris 通过 Multi-Catalog 对外部数据源做统一访问。在 MySQL 客户端里执行CREATE CATALOG paimon_catalog PROPERTIES ( type paimon, warehouse s3://my-bucket/warehouse/, s3.endpoint http://minio:9000, s3.access-key admin, s3.secret-key admin123 );创建成功后可以直接查询 Paimon 里的库和表SHOW DATABASES FROM paimon_catalog; SELECT COUNT(*) FROM paimon_catalog.dws_db.agent_trajectory WHERE dt 2025-01-20;这里有个细节查询时尽量带分区过滤条件。Doris 会把 WHERE 条件下推到 Paimon减少扫描量。我见过同事直接写不带 dt 的 COUNT(*) 去扫全表结果把对象存储的流量打满查询还特别慢。所以上线前一定要用 EXPLAIN 看执行计划确认过滤条件真的下推了。如果 Paimon 元数据存在 Hive Metastore也可以使用 Hive Catalog 方式挂载兼容性更稳。但那样会丢失部分 Paimon 的原生优化性能上不如直接 Paimon Catalog。4.3 用 Flink 把实时数据送进 Paimon 2.0实时链路我们统一用 Flink SQL。先在 Flink 里创建 Paimon CatalogCREATE CATALOG paimon WITH ( metastore filesystem, warehouse s3://my-bucket/warehouse/ );然后创建 Kafka SourceCREATE TABLE kafka_trajectory ( session_id STRING, event_id STRING, agent_id STRING, user_id STRING, action_type STRING, action_content STRING, observation STRING, ts TIMESTAMP(3), dt STRING ) WITH ( connector kafka, topic agent-trajectory, properties.bootstrap.servers kafka:9092, properties.group.id paimon_sink_group, format json, scan.startup.mode earliest-offset );再创建 Paimon 结果表。轨迹事件是追加型数据不需要主键直接用分区表CREATE TABLE paimon_db.agent_trajectory ( session_id STRING, event_id STRING, agent_id STRING, user_id STRING, action_type STRING, action_content STRING, observation STRING, ts TIMESTAMP(3), dt STRING ) PARTITIONED BY (dt);最后执行写入INSERT INTO paimon_db.agent_trajectory SELECT session_id, event_id, agent_id, user_id, action_type, action_content, observation, ts, DATE_FORMAT(ts, yyyy-MM-dd) FROM kafka_trajectory;Paimon 2.0 会在写入过程中自动做小文件合并但如果你把write-only设为 true则需要单独跑 Compaction 任务。我建议大多数场景保持默认让 Paimon 自动合并避免小文件失控。4.4 特征、反馈与训练样本如何回流在线特征查询我们不会每次都直接扫 Paimon 全表而是让 Doris 建物化视图。Doris 支持基于 External Catalog 的物化视图刷新CREATE MATERIALIZED VIEW mv_supplier_agg BUILD PERIODIC REFRESH EVERY 1 HOUR AS SELECT supplier_id, action_type, AVG(user_score) AS avg_score, COUNT(*) AS cnt, SUM(refund_flag) AS refund_cnt FROM paimon_catalog.dws_db.feedback GROUP BY supplier_id, action_type;这样 Doris 本地会存一份聚合结果Agent 在线查询时直接点查物化视图响应时间能到几十毫秒。至于反馈数据回流我们走 Kafka Flink 把人工评分写入 Paimon 的 feedback 主键表再用上面的物化视图定时刷新出来。这个方案有一个权衡物化视图刷新会有延迟。如果你的业务要求 1 分钟以内看到新反馈可以把刷新周期调小或者直接用 Doris 的同步物化视图。但刷新越频繁对底层存储的压力越大。我一般建议把反馈分成两层实时特征用 Doris 本地表承接 Kafka 流离线训练用 Paimon 全量数据两条链路在统一 schema 下并行。4.5 参数调优和冷热分层Paimon 2.0 有几个参数值得重点调。bucket决定主键表的哈希分桶数太小容易热点太大会产生大量小文件。我们一般按照 Flink 并行度来设置比如并行度 4 就用 4 个 bucket。snapshot.expire-limit和snapshot.time-retained用来控制历史快照保留时间如果不需要回放太老的数据可以调短节省存储。对象存储冷热分层也很有用。Agent 轨迹数据增长很快超过 30 天的冷数据可以放到低频存储。Paimon 支持把历史分区迁移到不同路径或者在对象存储层面配置生命周期策略。我们实践下来冷数据放在低频存储配合 Doris 查询时只访问最近分区成本能降一半以上。Doris 这边要关注外部表查询的内存。如果频繁查 Paimon建议给 Doris BE 预留足够的文件句柄和网络带宽。另外在线查询和离线分析最好走不同的 Doris 集群或不同计算组避免大查询把在线小查询拖死。5. 常见问题与排查技巧实录5.1 问题速查表下面这个表是我实际运维中遇到频率最高的问题可以直接当成排障清单用现象可能原因解决方案Doris 查 Paimon 很慢SQL 没有分区裁剪扫描全表检查执行计划强制 WHERE 带分区条件Paimon 小文件激增写入并发高自动合并滞后调整 bucket开启自动 Compaction主键表数据重复输入流存在重复事件未全局去重用 Flink 做 distinct或设置 sequence fieldAgent 在线查询延迟高Doris 并发打满或物化视图未命中增加物化视图分离在线和离线查询集群反馈数据对不上在线特征与离线训练口径不一致统一事件时间字段禁止混用到达时间挂载 Catalog 报错Doris 缺少 Paimon 依赖 jar替换 Doris 的 paimon-shaded jar 并重启 FE/BE5.2 版本兼容性踩坑Doris 通过外部函数读取 Paimon 表本质上是依赖 Paimon 的 Java API。如果你遇到的报错是UnsupportedFileFormat或者FileNotFoundException大概率是 Doris 内置的 Paimon 版本太低不认识新版 Paimon 写出的文件。解决办法是去 Doris 的fe/lib和be/lib目录把旧版 Paimon 相关 jar 替换成和你服务端匹配的版本。这里要特别小心不能只替换 FE 不替换 BE否则查询阶段仍可能报错。我们曾经只换了 FE结果还是查不到数据最后排查半天才发现 BE 没同步。还有一个兼容性经验如果公司已经用了 Hive Metastore并且不想额外维护 Paimon 的 filesystem catalog可以把 Paimon 的 metastore 配置成 hive。这样 Doris 用 Hive Catalog 也能访问虽然部分下推能力受限但胜在稳定适合快速验证架构。5.3 数据一致性怎么不翻车闭环系统最怕“闭环两头数字对不上”。我们的解决办法是统一时间口径。Agent 轨迹表里的时间一律用“事件发生时间”而不是“到达 Paimon 的时间”。因为离线任务只会看某个业务时间窗口如果混用了系统时间延迟数据可能污染窗口统计。另一个坑是主键表的删除语义。Agent 的反馈结果可能会被修正比如人工误操作后改分。如果直接把旧记录 DELETE 掉下游消费变更日志时可能会丢历史。我们建议在反馈表里加一个revised_flag字段采用“软更新”而不是物理删除保证历史可追溯。Doris 物化视图和时间窗口结合时也要小心。每小时刷新的物化视图如果刷新时间是 10:05但 Paimon 里面 10:00 的数据还没完全写入那这一批数据会少。我们通常把刷新时间和上游写入延迟错开并监控“最近一小时数据量”是否出现突降。5.4 验证闭环跑通的检查清单闭环有没有跑通不要只看一条链路能通就完事。我列一个检查清单从 Kafka 生产一条测试轨迹事件1 分钟内能在 Paimon 查到。Doris 通过 Catalog 能查到这条轨迹。Agent 在线查询能返回包含这条轨迹影响的特征。人工反馈事件回传后Paimon 的反馈表行数增加。Doris 物化视图刷新后聚合结果变化可见。离线训练任务能读取到新样本并输出一个新模型或新知识包。新模型发布后Agent 下一次回答命中了新数据。我习惯在 Doris 里建一张“闭环状态表”每隔十分钟统计 Kafka 的消费 offset、Paimon 各分区行数、物化视图刷新时间。只要这张表的数据持续增长我就放心闭环还在转。6. 一些落地体会最后说点个人体会。Agentic AI 的数据闭环难点不在选某个开源组件有多新而在于你能不能把“反馈”这件事定义清楚。很多团队连“什么样的行为算好行为”都没定义就开始搭湖仓结果数据全存下来了却永远不知道哪部分应该回灌给模型。Apache Doris Paimon 2.0 的组合让我们真正做到了一边保留全量可回放的 Agent 轨迹一边给在线决策提供稳定的低延迟计算能力。如果你也在做 Agentic AI我建议不要一上来就上一整套复杂的湖仓平台先用这套组合把一两个核心场景的闭环跑通再逐步扩大。最后分享一个小诀窍用 Doris 的审计日志配合 Agent 请求日志可以看到每个 Agent 调用了哪些特征表、哪些 SQL 最频繁。这些信息反过来能指导你调整物化视图和 Paimon 分区策略。数据闭环不是静态架构它自己也需要持续迭代。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门