最近在帮团队搭一个 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 执行动作后产生事件;事件经过消息队列进入 Paimon;Paimon 的增量数据再回流到 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 天的订单量、好评率、平均响应时间;然后调外部服务获取实时价格;最后综合这些信息给出推荐。
用户选择了其中一家,或者给了差评。这些行为进入 Kafka,Flink 消费后写入 Paimon 的轨迹表;离线任务扫描新的反馈,更新每家公司的评分;更新后的评分写回 Paimon 或 Doris 的特征表。第二天,另一个用户再问相同的问题时,Agent 会拿到新的评分数据,从而改变推荐结果。
这就是闭环真正合上的位置:不是某一条数据管道跑通了,而是“分析结果能反过来影响下一次决策”。Doris 和 Paimon 分别承载了这个环的两端。
4. 实操:打通 Doris 与 Paimon 的关键步骤
4.1 环境和版本准备
先说版本。我们用的是 Doris 2.1.4,Paimon 2.0.1,Flink 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 Catalog
Doris 通过 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 Catalog:
CREATE CATALOG paimon WITH ( 'metastore' = 'filesystem', 'warehouse' = 's3://my-bucket/warehouse/' );然后创建 Kafka Source:
CREATE 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 field |
| Agent 在线查询延迟高 | Doris 并发打满,或物化视图未命中 | 增加物化视图,分离在线和离线查询集群 |
| 反馈数据对不上 | 在线特征与离线训练口径不一致 | 统一事件时间字段,禁止混用到达时间 |
| 挂载 Catalog 报错 | Doris 缺少 Paimon 依赖 jar | 替换 Doris 的 paimon-shaded jar 并重启 FE/BE |
5.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 分区策略。数据闭环不是静态架构,它自己也需要持续迭代。