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

资讯详情

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

基于SelectDB构建PB级Agent可观测平台:架构设计与实战

基于SelectDB构建PB级Agent可观测平台:架构设计与实战 1. 项目概述当Agent智能体遇上PB级数据洪流最近和几个做AI应用的朋友聊天大家普遍都在头疼同一个问题当你的智能体Agent系统真正跑起来每天处理成千上万的用户交互产生海量的日志、追踪数据和性能指标时该怎么管怎么查怎么分析这已经不是简单的“看日志”能解决的了。一个复杂的Agent工作流可能涉及多轮对话、工具调用、外部API请求、内部状态流转每一次交互都是一条长长的“轨迹”。日积月累数据量轻松突破TB级奔着PBPetabyte千万亿字节级别就去了。传统的日志系统或者时序数据库在应对这种既要高吞吐写入又要支持复杂即席查询Ad-hoc Query的场景时往往力不从心。这正是“阶跃星辰基于 SelectDB 构建 PB 级 Agent 可观测平台”这个项目要解决的核心痛点。简单来说它构建了一个能“看清”超大规模AI智能体系统内部每一个动作、每一次思考、每一处瓶颈的“超级望远镜”。这个平台不仅要能吞下海量数据还要能让研发、运维乃至业务同学快速、灵活地从这些数据中挖掘出价值比如某个Agent的响应为什么变慢了某个工具调用的成功率为何下降用户与Agent的交互模式发生了什么变化SelectDB 在这里扮演了核心“数据引擎”的角色。它不是一个简单的存储库而是一个能够对PB级数据进行实时分析与交互式查询的现代化分析型数据库。这个项目的成功意味着大规模AI应用的可观测性从“奢侈品”变成了“标配”为AI系统的稳定运营、持续优化和商业洞察提供了坚实的数据地基。无论你是AI应用开发者、SRE工程师还是数据平台架构师理解这套方案的思路与实现都极具参考价值。2. 平台整体架构设计与核心思路拆解构建一个面向PB级Agent数据的可观测平台绝非将数据往一个强大的数据库里一丢了之。它需要一整套从数据采集、传输、处理到存储、分析、可视化的完整架构设计。其核心思路可以概括为“全链路埋点、流批一体入湖、统一分析服务”。2.1 为什么是“可观测性”而不仅仅是“监控”首先需要厘清一个概念监控Monitoring和可观测性Observability有本质区别。传统监控是针对已知的、预设的指标如CPU使用率、错误率进行阈值告警。而可观测性面对的是未知的、复杂系统的内部状态它通过收集日志Logs、指标Metrics和追踪Traces这三大支柱数据允许我们提出任意的新问题并通过探索数据来找到答案。对于Agent系统而言可观测性尤为重要。你无法预知用户会问出什么千奇百怪的问题也无法预知Agent在复杂决策链中会在哪个环节“卡住”。因此平台必须能记录下完整的“轨迹”日志记录每个步骤的详细信息如接收的用户输入、调用的模型名称、生成的思考过程、执行工具的参数和结果。指标聚合的性能数据如每秒查询量QPS、平均响应延迟、各阶段耗时思考、行动、生成的P99分位值、工具调用成功率。追踪一次完整的用户会话Session在系统内部跨越多个服务、函数和外部调用的全链路调用关系与耗时通常用唯一的TraceID串联。这个项目的架构正是围绕高效收集、存储和查询这三大类数据而展开的。2.2 技术选型为什么是SelectDB面对PB级、需要支持高并发即席查询的场景技术选型是成败关键。常见的选项有Elasticsearch强在全文检索和简单的聚合但对于多表关联、复杂SQL分析能力较弱在PB级规模下存储和计算成本较高。时序数据库如InfluxDB, TDengine为指标数据优化写入和压缩效率高但对日志和追踪这类非规整、需要复杂查询的数据模型支持不友好。数据湖如Iceberg/Hudi 计算引擎如Trino/Flink存算分离灵活性极高但通常查询延迟在秒到分钟级难以满足交互式诊断的亚秒级响应需求。现代分析型数据库如SelectDB, ClickHouse专为大规模数据分析设计兼容标准SQL支持高吞吐写入和低延迟复杂查询。SelectDB在此脱颖而出主要基于以下几点考量极致的分析性能基于MPP大规模并行处理架构和列式存储对聚合、过滤、多表JOIN等分析查询有天然优势。这对于分析Agent的会话路径、统计工具使用频率、进行漏斗分析等场景至关重要。对标准SQL的完整支持降低了使用门槛数据分析师和研发可以直接用熟悉的SQL进行探索无需学习新的查询语言如ES的DSL。这对于构建统一的数据分析平台至关重要。良好的生态集成支持作为Flink、Kafka等流处理引擎的Sink便于实现流批一体的数据实时入库。同时支持多种数据湖表格式如Iceberg的联邦查询为未来接入更原始的数据提供了扩展性。成本与效率的平衡通过智能索引、物化视图、数据分区与分桶等机制在PB级数据量下依然能保持高效的查询性能同时控制存储成本。因此选择SelectDB作为核心分析引擎是在性能、灵活性、易用性和成本之间做出的一个最优平衡。2.3 整体架构蓝图一个典型的平台架构会分为四层数据采集层在Agent框架的各个关键节点植入轻量级SDK以非侵入或低侵入的方式收集日志、指标和追踪数据。数据格式通常统一为JSON通过异步方式发送到消息队列如Kafka避免阻塞主业务逻辑。数据处理与接入层使用流处理引擎如Apache Flink消费Kafka中的数据进行必要的清洗、格式化、富化例如补充用户标签、会话信息和轻量聚合。处理后的数据流通过SelectDB提供的原生连接器如selectdb-sink或通过INSERT INTO SELECT的方式实时写入SelectDB的明细表中。存储与分析层这是SelectDB的主场。根据数据类型和查询模式设计不同的表结构明细事件表存储每一条原始的Agent轨迹事件采用分区按天/小时和分桶按TraceID或UserID哈希策略支持最细粒度的回溯与诊断。聚合指标表通过物化视图Materialized View或定时任务将明细数据预聚合为分钟级、小时级的性能指标表用于仪表盘和告警极大提升高频查询速度。维度表存储相对静态的信息如Agent版本号、工具定义、用户分区等用于关联分析。应用与可视化层基于SelectDB提供统一的SQL查询接口。上游可以连接BI工具如Metabase、Superset用于制作面向业务和运营的仪表盘如每日活跃Agent、平均会话长度、用户满意度趋势。自定义可观测平台开发前端界面提供轨迹追踪树Trace Tree可视化、会话回放、性能对比分析等高级功能。告警系统通过定时查询或监听物化视图对异常指标如错误率飙升、延迟陡增触发告警。注意架构设计初期就必须考虑数据生命周期管理TTL。PB级数据不可能永久保存。通常策略是原始明细数据保留7-30天用于问题诊断聚合后的指标数据保留1-2年用于趋势分析。SelectDB可以通过PARTITION和DROP PARTITION语句方便地实现这一点。3. 核心细节解析数据模型设计与查询优化平台的能力上限很大程度上由数据模型决定。设计一个既能满足灵活查询又能保证PB级规模下高效运行的表结构是核心中的核心。3.1 Agent可观测数据模型设计我们需要设计一张或多张核心的明细表。这里给出一个高度简化的核心事件表agent_events的设计示例它试图用一种相对通用的格式记录Agent执行过程中的关键步骤CREATE TABLE agent_events ( event_id BIGINT COMMENT 事件唯一ID, trace_id VARCHAR(256) NOT NULL COMMENT 全链路追踪ID贯穿一次会话, span_id VARCHAR(256) COMMENT 当前跨度ID, parent_span_id VARCHAR(256) COMMENT 父跨度ID, event_time DATETIMEV2(3) NOT NULL COMMENT 事件发生时间精确到毫秒, agent_id VARCHAR(128) COMMENT Agent实例或类型标识, session_id VARCHAR(256) NOT NULL COMMENT 用户会话ID, user_id VARCHAR(128) COMMENT 用户标识, event_type VARCHAR(64) NOT NULL COMMENT 事件类型如user_input, llm_invoke, tool_call, tool_result, final_output, error, event_stage VARCHAR(64) COMMENT 事件阶段如thinking, action, observation, content STRING COMMENT 事件内容JSON格式存储详细信息, model_name VARCHAR(128) COMMENT 调用的LLM模型名称, tool_name VARCHAR(128) COMMENT 调用的工具名称, duration_ms INT COMMENT 该步骤耗时毫秒, status_code VARCHAR(32) COMMENT 状态码如SUCCESS, FAILED, TIMEOUT, error_message STRING COMMENT 错误信息, tags MAPVARCHAR, VARCHAR COMMENT 自定义标签用于灵活过滤和分组, __partition_date DATE NOT NULL COMMENT 分区字段按天分区 ) ENGINEOLAP DUPLICATE KEY(trace_id, event_time, event_type) -- 重复数据模型适合明细分析 PARTITION BY RANGE(__partition_date)() -- 动态分区后续按天添加 DISTRIBUTED BY HASH(trace_id) BUCKETS 32 -- 按trace_id哈希分桶保证同一会话数据在相同节点 PROPERTIES ( replication_num 3, storage_medium SSD, dynamic_partition.enable true, dynamic_partition.time_unit DAY, dynamic_partition.start -7, -- 保留最近7天分区 dynamic_partition.end 3, dynamic_partition.prefix p );设计要点解析trace_id作为核心键DISTRIBUTED BY HASH(trace_id)是关键决策。这确保了同一次用户会话Trace的所有事件Spans有很大概率分布在同一个数据分片Bucket内。当我们需要重构整个会话轨迹时可以极大减少跨节点的数据拉取提升查询性能。event_time高精度使用DATETIMEV2(3)存储毫秒级时间戳对于分析微秒级延迟分布至关重要。content字段使用STRING/JSON类型Agent内部状态、LLM的输入输出、工具调用的具体参数和结果可能结构多变且嵌套很深。将其存储为JSON字符串既能保持灵活性又能利用SelectDB的JSON函数进行部分查询。但需注意对JSON字段内属性的查询效率通常低于扁平化的列。tagsMAP类型这是一个非常实用的设计。用于存储各种动态的、可能后期增加的维度信息如实验分组、流量来源、客户端版本等。MAP类型支持高效的key-value查询和过滤。分区与分桶按日期分区是管理海量数据生命周期的标准做法。按trace_id哈希分桶则是针对查询模式做的优化。3.2 物化视图与聚合表加速查询对agent_events明细表的直接查询在数据量巨大时即使对于SelectDB也可能较慢尤其是那些需要扫描大量数据的聚合查询如“今天所有会话的平均响应时间”。此时物化视图Materialized View是性能加速的利器。例如我们可以创建一个每分钟聚合一次的核心性能指标物化视图-- 基于明细表创建物化视图自动预聚合 CREATE MATERIALIZED VIEW mv_agent_perf_minute BUILD IMMEDIATE -- 立即构建 REFRESH ASYNC -- 异步刷新 KEY(agg_minute, agent_id, status_code) PARTITION BY (agg_minute) DISTRIBUTED BY HASH(agent_id) BUCKETS 10 AS SELECT DATE_TRUNC(minute, event_time) AS agg_minute, agent_id, status_code, COUNT(*) AS event_count, SUM(duration_ms) AS total_duration, COUNT(DISTINCT session_id) AS active_sessions, COUNT(DISTINCT trace_id) AS request_count FROM agent_events WHERE event_type IN (llm_invoke, tool_call, final_output) GROUP BY agg_minute, agent_id, status_code;创建此物化视图后当用户查询每分钟的性能指标时SelectDB的优化器会自动将查询路由到这个小得多的、已预聚合的物化视图上查询速度可能提升数百倍。实操心得物化视图虽好但不宜过多。每个物化视图都是存储和计算开销。优先为查询最频繁、数据量最大、且聚合模式固定的查询创建。对于更复杂的、多维度的即席查询可能还需要配合创建专门的聚合表通过定时任务如Flink Job来更新。3.3 索引策略Bloom Filter与倒排索引对于agent_events这种宽表高效的过滤是快速查询的前提。除了利用分区和分桶进行粗粒度剪枝我们还需要在列级别建立索引。前缀索引SelectDB默认使用前缀索引基于DUPLICATE/UNIQUE KEY的前36字节。在我们的设计中trace_id和event_time在最前面对于按trace_id或时间范围查询非常有效。Bloom Filter索引对于event_type,agent_id,status_code这类高基数列值种类多的等值查询创建Bloom Filter索引可以快速判断数据块中是否包含目标值避免不必要的扫描。ALTER TABLE agent_events ADD INDEX bf_idx_event_type (event_type) USING BLOOM_FILTER; ALTER TABLE agent_events ADD INDEX bf_idx_agent (agent_id) USING BLOOM_FILTER;倒排索引从2.0版本开始支持Beta对于content或error_message这类长文本字段如果需要经常进行关键词搜索如查找包含“超时”错误的事件可以尝试使用倒排索引来加速文本检索。索引创建原则按需创建持续监控。通过分析慢查询日志找到频繁用于WHERE条件且过滤效果好的列为其添加合适的索引。避免盲目添加因为索引会占用额外空间并影响写入速度。4. 实操过程从数据接入到可视化分析有了清晰的数据模型和表结构下一步就是让数据流动起来并最终产生价值。这个过程涉及多个系统的协同。4.1 数据采集与实时写入流水线假设我们使用Apache Flink作为流处理引擎负责消费Kafka中的原始事件并进行ETL后写入SelectDB。步骤一定义Flink SourceKafka配置Flink从Kafka主题agent-raw-events中读取JSON格式的原始数据。步骤二定义Flink Transformation在Flink作业中对数据进行解析、清洗和富化。例如解析JSON提取字段。补充trace_id、session_id如果上游未生成。将嵌套的复杂对象序列化成字符串存入content字段。计算duration_ms如果上下游事件有时间戳。添加一些业务标签到tagsMAP中。步骤三定义Flink SinkSelectDB使用SelectDB官方提供的selectdb-sink连接器或者使用通用的JDBC Sink。以下是使用selectdb-sink的简化配置示例// Flink SQL 方式示例 String sinkDDL CREATE TABLE selectdb_sink (\n ... 字段定义需与SelectDB表对齐 ...\n ) WITH (\n connector selectdb,\n table.identifier db_name.agent_events,\n username user,\n password pass,\n sink.buffer-flush.max-rows 50000,\n sink.buffer-flush.interval 10s,\n sink.properties.format json,\n sink.properties.read_json_by_line true,\n sink.max-retries 3\n ); tableEnv.executeSql(sinkDDL); // 将处理后的数据流插入到sink表 tableEnv.executeSql(INSERT INTO selectdb_sink SELECT ... FROM transformed_stream);关键配置解析sink.buffer-flush.*控制写入的批处理行为。max-rows和interval共同决定了数据攒批的条件达到任一条件即触发一次批量写入。这是平衡写入吞吐和实时性的关键参数。对于日志类数据可以适当调大如50万行/30秒以获得更高的吞吐。sink.max-retries网络抖动或SelectDB短暂不可用时重试机制能保证数据不丢失。sink.properties.format设置为json与SelectDB表的Stream LoadJSON导入模式匹配写入效率高。踩坑记录在早期测试中我们曾使用单条INSERT语句的JDBC Sink写入性能极差且对SelectDB前端节点FE压力巨大。切换到使用selectdb-sink底层基于Stream Load后写入吞吐量从每秒几百条提升到每秒数十万条这是构建实时可观测平台的基石。务必使用批量的、面向分析数据库优化的写入方式。4.2 构建交互式可观测查询数据入库后我们就可以通过SQL这个强大的工具进行各种分析了。以下是一些典型场景的查询示例场景一故障诊断 - 追溯一个失败会话的完整轨迹SELECT event_time, event_type, event_stage, tool_name, model_name, duration_ms, status_code, error_message, content -- 查看详细内容 FROM agent_events WHERE trace_id your_failed_trace_id_here ORDER BY event_time ASC;这个查询利用了trace_id的分桶优化能快速拉取出一次会话的所有事件按时间排序后就能像看剧本一样复盘Agent的整个“思考-行动”过程精准定位是在调用哪个工具、询问哪个模型时出了错。场景二性能分析 - 统计今日各Agent的平均响应时间及P99延迟SELECT agent_id, COUNT(*) as total_requests, AVG(duration_ms) as avg_latency, PERCENTILE_APPROX(duration_ms, 0.99) as p99_latency -- SelectDB的近似百分位函数性能极高 FROM agent_events WHERE __partition_date 2023-10-27 AND event_type final_output -- 以最终输出事件作为一次请求的结束 AND status_code SUCCESS GROUP BY agent_id ORDER BY p99_latency DESC;这个查询能快速找出性能瓶颈所在的Agent。PERCENTILE_APPROX函数对于分析长尾延迟P99, P999非常有用且计算效率远高于精确百分位。场景三用量与成本分析 - 分析不同LLM模型的调用量与平均token消耗假设我们在content字段的JSON中记录了usage信息。SELECT model_name, COUNT(*) as invoke_count, AVG(CAST(JSON_EXTRACT(content, $.usage.prompt_tokens) AS INT)) as avg_prompt_tokens, AVG(CAST(JSON_EXTRACT(content, $.usage.completion_tokens) AS INT)) as avg_completion_tokens FROM agent_events WHERE __partition_date 2023-10-20 AND event_type llm_invoke AND model_name IS NOT NULL GROUP BY model_name;这个查询直接关联了可观测数据和成本对于优化模型调用策略、控制预算至关重要。场景四会话路径分析 - 找出用户最常使用的工具组合序列模式这需要用到窗口函数进行会话内排序和相邻事件分析相对复杂但SelectDB的SQL能力足以支持。WITH session_events AS ( SELECT session_id, event_type, tool_name, LAG(tool_name) OVER (PARTITION BY session_id ORDER BY event_time) as prev_tool FROM agent_events WHERE __partition_date 2023-10-27 AND event_type tool_call ) SELECT CONCAT(prev_tool, - , tool_name) as tool_sequence, COUNT(*) as frequency FROM session_events WHERE prev_tool IS NOT NULL GROUP BY tool_sequence ORDER BY frequency DESC LIMIT 10;这个查询结果可以直观展示用户与Agent交互中最常见的“工作流”为产品优化和工具链设计提供洞察。4.3 可视化与告警集成数据分析的结果需要以更直观的方式呈现。BI集成将SelectDB作为数据源连接到Metabase。可以轻松创建仪表盘例如实时监控大盘展示当前QPS、错误率、平均延迟等核心指标数据来自物化视图mv_agent_perf_minute。会话分析看板展示会话深度分布、常用工具排行榜、用户地域分布等。成本分析看板展示各模型每日token消耗趋势、成本占比。自定义可观测平台对于更专业的可观测需求如Trace树可视化需要后端服务查询SelectDB的agent_events表按trace_id和span_id、parent_span_id的关系在内存中构建树形结构然后前端使用类似ECharts或AntV G6的图形库进行渲染。告警可以编写定时脚本查询SelectDB中的聚合指标当错误率超过5%或P99延迟超过设定阈值时通过Webhook触发告警如发送到钉钉、飞书或PagerDuty。更优雅的方式是使用Grafana的告警功能直接对接SelectDB数据源。5. 常见问题与性能调优实战录在PB级规模下运行这样一个平台挑战无处不在。以下是我们实践中遇到的一些典型问题及解决方案。5.1 写入性能瓶颈与优化问题现象数据写入速度跟不上采集速度Flink Checkpoint时间变长甚至出现背压Backpressure。排查与解决检查SelectDB Sink配置首先确认selectdb-sink的buffer-flush参数是否合理。如果设置得太小如1万行/1秒会导致过于频繁的小批量写入增加网络往返和SelectDB的导入事务开销。可以适当调大例如设置为‘sink.buffer-flush.max-rows’ ‘100000’和‘sink.buffer-flush.interval’ ‘20s’在内存允许的情况下积累更大批次再写入。观察SelectDB BE节点负载通过SelectDB的监控指标或SHOW PROC ‘/backends’命令查看各个后端节点Backend的CPU、内存、磁盘IO和网络流量。如果某个BE负载明显偏高可能是数据分桶不均匀导致。需要调整建表时的DISTRIBUTED BY HASH的列或者增加BUCKETS数量使数据更均匀地分布。检查表结构是否合理过宽的表数百列或包含超大字符串列如存储整个对话历史会影响写入吞吐。评估是否可以将一些不常查询的大字段剥离到另一张表通过event_id关联。升级SelectDB版本并利用新特性新版本可能在写入路径上有优化。例如确认是否使用了高效的Stream Load协议并开启了写入压缩。5.2 查询响应慢与优化问题现象面向明细表的即席查询特别是涉及多天数据扫描或复杂条件过滤的查询响应时间超过10秒无法满足交互式分析需求。排查与解决利用EXPLAIN分析查询计划在查询前加上EXPLAIN查看执行计划。重点关注分区裁剪Partition Pruning是否有效利用了__partition_date等分区字段只扫描了必要的分区。分桶裁剪Bucket Pruning如果查询条件包含分桶键trace_id的等值条件是否只扫描了对应的分桶。索引命中是否使用了Bloom Filter等索引进行数据块过滤。优化SQL写法避免使用SELECT *只查询需要的列减少数据传输和序列化开销。将过滤条件尽可能提前尤其是在子查询中尽早过滤掉无关数据。谨慎使用LIKE ‘%...%’全模糊匹配对字符串列的全模糊匹配无法利用索引性能极差。如果必须使用考虑对该列建立倒排索引。注意OR条件的性能多个OR条件可能导致索引失效。尝试改写为UNION ALL或者使用更高效的过滤结构。利用物化视图这是解决聚合查询慢的最有效手段。分析慢查询日志找出频繁且模式固定的聚合查询为其创建物化视图。调整数据模型如果发现某些JSON字段中的属性被频繁查询可以考虑将其提取出来作为单独的列建立索引从而获得最佳的查询性能。5.3 存储成本膨胀问题问题现象数据量增长过快存储成本超出预期。排查与解决实施严格的数据生命周期策略这是最直接有效的方法。通过ALTER TABLE agent_events DROP PARTITION p20231001;删除过期分区。结合动态分区属性实现自动管理。评估数据压缩率SelectDB默认使用LZ4压缩。对于文本日志类数据可以尝试使用ZSTD压缩算法建表时设置“compression” “ZSTD”通常能获得更高的压缩比但会略微增加CPU开销。冷热数据分层存储将近期热数据放在SSD上以保证性能将历史冷数据自动转存到更廉价的S3/HDFS对象存储上。SelectDB支持冷热数据分离配置。数据降粒度归档对于超过一定时间如30天的明细数据可以定期运行归档任务将其聚合为小时级或天级的概要数据然后删除原始明细。只保留必要的细节用于长期趋势分析。5.4 数据一致性与准确性保障问题现象监控发现写入的数据量与应用端发送的量对不上或者查询结果出现重复。排查与解决确保端到端至少一次At-Least-Once语义在Flink到SelectDB的链路中启用Flink Checkpoint和SelectDB Sink的两阶段提交2PC如果支持或幂等写入可以保证在故障恢复时数据不丢失。但“至少一次”可能带来重复。处理重复数据SelectDB的DUPLICATE KEY模型允许重复数据。如果业务不能接受重复需要在应用层或Flink层实现幂等如根据event_id去重或者在SelectDB中通过聚合查询来DISTINCT。监控数据延迟在Flink作业中增加一个指标记录事件产生时间event_time和写入SelectDB时间process_time的差值即端到端延迟。通过监控这个差值可以及时发现数据处理链路的堵塞。建立数据质量校验任务定期运行一些基准SQL检查数据是否连续、关键字段的非空率是否正常、数值指标是否在合理范围内。将校验结果纳入监控告警。构建这样一个平台是一个持续迭代和优化的过程。从最初的原型到支撑PB级数据我们经历了多次架构调整、参数调优和问题攻坚。最大的体会是可观测性建设必须与业务发展同步甚至超前规划。当你的Agent智能体开始承担核心业务流量时一个稳定、高效、洞察力强的可观测平台就是你系统中最可靠的“护航员”。它不仅能帮你快速灭火更能通过数据驱动指引你优化Agent的“智力”与“体力”最终提升用户体验和业务价值。
返回列表