AgentOps Exporter:从 Supabase (Postgres) 到 ClickHouse 的 OpenTelemetry 数据迁移与双写通道
2026/9/17 22:18:54 网站建设 项目流程

AgentOps Exporter:从 Supabase (Postgres) 到 ClickHouse 的 OpenTelemetry 数据迁移与双写通道

【免费下载链接】agentopsPython SDK for AI agent monitoring, LLM cost tracking, benchmarking, and more. Integrates with most LLMs and agent frameworks including CrewAI, Agno, OpenAI Agents SDK, Langchain, Autogen, AG2, and CamelAI项目地址: https://gitcode.com/GitHub_Trending/ag/agentops

AgentOps 的后端正处于从传统关系型存储(Supabase/Postgres)向 OpenTelemetry + ClickHouse 架构迁移的过程中,app/api/agentops/exporter/目录就是这次数据搬迁的核心通道。本文基于该目录的 README 及其三个核心源码文件(export.pyprocessor.pymodels.py),完整讲解 Exporter 的两类使用方式:作为遗留 API 的实时双写处理器、作为可断点续跑的批量迁移工具,以及底层 "Session/Agent/Event → Trace/Span" 的 OTEL 语义映射规则,帮助读者理解 AgentOps 如何将 92.5 GB 级别的存量 Agent 会话数据转换为标准的 OTEL Traces 格式。

Exporter 的定位与双重角色

根据 Exporter README 的原始描述,该服务承担两个职责:

  1. 一次性迁移(One-Time Migration):一个可按需执行的脚本,把 Postgres 存量 schema 中的所有数据迁入 ClickHouse;
  2. 并行处理器(Parallel Processor):数据仍按原有 schema 写入原有服务端(Supabase),同时一条并行数据流把 OTEL 格式的数据写入 ClickHouse。

processor.py的模块 docstring 也印证了这一双重定位:该模块"既是 Supabase 与 ClickHouse 之间数据转换的内部 API,也是一个把 Supabase 数据批量导出到 ClickHouse 的 CLI 工具"(见 processor.py)。

README 中还记录了当时的迁移背景数据:存量数据约92.5 GB,以及一份 v3 API(仓库中实际前缀为/v2)的端点改名草图——session → traceevents → spanagent → span (that holds other spans),并明确了 "every event belongs to an agent" 的层级关系。这套命名草图在models.py中得到了完整实现。

双写模式:v2 API 端点如何触发 Exporter

实时双写发生在 FastAPI 路由层。v2.py 在文件顶部导入from agentops.exporter import export,并在每个数据落库端点中紧跟 Supabase 写入调用对应的 export 函数:

  • POST /v2/sessionsPOST /v2/create_sessionawait supabase.table("sessions").insert(session).execute()之后紧跟await export.create_session(session)(见 v2.py);
  • POST /v2/update_session:先supabase.table("sessions").update(...)await export.update_session(session)
  • POST /v2/create_agentagents表 upsert 后调用export.create_agent(agent)
  • POST /v2/create_events:按event_type把事件分流为llms/tools/errors/actions四类,分别写入对应表后,逐条调用export.create_llm_event/create_tool_event/create_error_event/create_action_event(见 v2.py);
  • POST /v2/update_events:对每类事件执行 Supabase update 后调用对应的export.update_*函数。

也就是说,Supabase 写入是主路径,ClickHouse 写入是并行影子流,两条流共享同一份经过event_handlers处理后的数据结构。

映射模型:Session 变 Trace,事件变 Span

models.py 用 Pydantic 定义了源端六张表到 OTEL 对象的映射。EXPORT_TABLES_MODELS(见 processor.py)给出了完整的表-模型对照:

Supabase 表Pydantic 模型OTEL 目标形态SpanKind
sessionsSessionTrace 的根 Span(span_id == trace_id == session.idSERVER
agentsAgent父级 Span,parent_span_id = session_idINTERNAL
actionsActionEvent子 Span,parent_span_id = agent_idINTERNAL
llmsLLMEvent子 Span,parent_span_id = agent_idCLIENT
toolsToolEvent子 Span,parent_span_id = agent_idINTERNAL
errorsErrorEvent点事件 Span,parent_span_id = session_idStatusCode = ERRORINTERNAL

这与 README "Noise" 部分的草图一一对应:/create_session/update_session废弃 "session" 一词转为 trace;/create_events/update_events将 events 转为 span;/create_agent将 agent 转为"持有其他 span 的 span"。

几个关键实现细节:

  • ID 复用export.py的 docstring 点明核心规则——"Traces and their root spans share the same ID which comes from the Session ID"(见 export.py)。因此TraceId就是session_id,span 的SpanId就是各表的主键 UUID,无需生成新 ID,天然满足幂等对齐;
  • 时间单位转换:ClickHouse 要求纳秒/DateTime64,datetime_to_ns()datetime转为纳秒,把整型时间戳视为毫秒再乘 1e6(见 models.py);
  • LLM 事件采用 GenAI 语义约定LLMEvent.to_span()生成gen_ai.prompt.{i}.contentgen_ai.completion.{i}.finish_reasongen_ai.usage.prompt_tokens等属性,同时保留llm.modelllm.costllm.total_tokens等 AgentOps 私有属性,并标注了 GenAI 语义约定的来源(OpenTelemetry 官方规范),见 models.py;
  • Span 到 ClickHouse 行的序列化Span.to_clickhouse_dict()输出的列名与DESCRIBE otel_traces完全对齐(TimestampTraceIdSpanIdParentSpanIdSpanNameSpanKindResourceAttributesEvents.Timestamp等),并在写入前强制注入ResourceAttributes["agentops.project.id"] = project_id(见 models.py)。

这个agentops.project.id注入不是普通属性:ClickHouse 侧的otel_traces表通过MATERIALIZED列直接派生project_id,见 0000_init.sql:

`project_id` String MATERIALIZED ResourceAttributes['agentops.project.id'], ... PARTITION BY toYYYYMM(Timestamp) ORDER BY (project_id, Timestamp)

otel_traces按月分区、按(project_id, Timestamp)排序并带有bloom_filter索引(idx_trace_ididx_span_ididx_project_id)。因此 Exporter 正确写入agentops.project.id是多租户检索(前端/v4/traces?project_id=...)能正常工作的前提。

实时写入路径:export.py 的调用链

export.py 对外暴露一组异步函数,供 API 层同步调用。以创建 LLM 事件为例,调用链是:

  1. 校验data中必须携带session_idagent_id
  2. _get_project_id(session_id)通过 Supabase 客户端反查sessions表拿到project_id
  3. 构造 Pydantic 模型(如LLMEvent(**data)),调用to_span(trace_id=session_id, parent_span_id=agent_id, project_id=project_id)生成Span
  4. clickhouse_create_span(span)将其序列化为 ClickHouse 行并插入otel_traces表。

更新路径(update_*)则走clickhouse_update_span,其策略比较特殊:由于 ClickHouse 的 Map/Array 列原地 UPDATE 序列化复杂,实现选择ALTER TABLE ... DELETE再合并重插("删除+重插的传播延迟与 UPDATE 相当",见 processor.py)。重插前会用merge_dicts_recursive把库中已有行与新数据做递归字典合并,避免丢失未更新的列。

clickhouse_create内部还有一层值转义:clickhouse_escape_value处理None(转NULL)、list、dict、int,再配合clickhouse_driverescape_param(见 processor.py);DRY_RUN = True时只打印不写库,便于灰度验证。

批量迁移路径:processor.py 的分页、并行与断点

processor.py既是被export.py引用的内部库,也是可独立运行的 CLI 入口(if __name__ == "__main__"load_dotenv('.env', override=True)asyncio.run(main()),见 processor.py)。核心参数都集中定义在文件头部:

SUPABASE_MIN_POOL_SIZE = 12 # Supabase 连接池下限 SUPABASE_MAX_POOL_SIZE = 24 # 上限 DRY_RUN = False # True 时不写 ClickHouse PAGE_NUMBER = 0 MAX_PAGES = None # None 表示处理全部页 TIMEOUT = 240 # 单条 trace 写入超时(秒) BATCH_ROW_COUNT = 1000 # 每页从 Supabase 拉取的行数 MAX_CONCURRENT = 42 # 最大并发写入任务数 PARALLEL_READS = 21 # 并行读取 session 的批大小

批量流程(main(),见 processor.py):

  1. init_files()初始化两个持久化文件:dropped_records.csv(记录迁移中因异常被丢弃的记录)与last_session_id.txt(断点检查点);
  2. count_session_rows()统计剩余 sessions 行数,total_pages = row_count // BATCH_ROW_COUNT
  3. 逐页调用process_page(offset, limit)get_sessionsid > last_session_id ORDER BY id ASC LIMIT/OFFSET分页拉取(UUID 按字典序排),再以PARALLEL_READS=21为批并行执行get_session_as_trace
  4. process_page维护一个最多MAX_CONCURRENT=42的待写入任务池,每条 trace 通过write_trace_with_timeoutasyncio.wait_for(..., timeout=240))写入,超时或异常的记录落到dropped_records.csv而不是中断整体迁移;
  5. get_session_as_trace中每成功处理一个 session 就调用write_last_session_id把当前session_id回写到last_session_id.txt(seek + truncate + write),因此任务中断后重跑可从断点继续,不会重复导出已完成的 session;
  6. MAX_PAGES非空时处理到指定页数后提前退出——这正是 README "Next Steps" 中SELECT ... LIMIT 0, 10000小批量试跑策略的落地手段。

Supabase 访问层有一个值得注意的设计:SupabaseExporter通过元类SupabaseExporterMeta支持supabase[Session]这种"按模型取表"的语法糖,fetchall模板串里固定用{table_name}占位(见 processor.py),连接串由SUPABASE_HOST/PORT/DATABASE/USER/PASSWORD环境变量拼成postgresql://DSN。

此外还有一个针对历史脏数据的修复工具:get_v2_sourced_rowsotel_2.otel_traces中筛出携带session.id属性的行(即早期由 v2 exporter 写入的数据),再通过assign_correct_project_id执行ALTER TABLE ... UPDATE ResourceAttributes = mapUpdate(...)逐行修正agentops.project.id,配合get_project_id_for_session_id的 1 万条 LRU 缓存查 Supabase 补全正确项目归属。

README 中"Next Steps"小样本验证流程

原 README 记录了一条完整的低风险验证路径,结合源码可以落地为:

  1. 启动一个 dev 环境的 Supabase Postgres 实例,并准备一个一次性的(throwaway)ClickHouse 实例;
  2. 从生产库选取一段样本(思路为SELECT ... LIMIT 0, 10000,对应批量参数BATCH_ROW_COUNT=1000+MAX_PAGES限制页数),写入 dev Supabase;
  3. 启动 exporter 处理实例(python processor.py直接运行),通过DRY_RUN先干跑确认输出结构,再正式写入;
  4. dropped_records.csv核对失败记录,用last_session_id.txt确认迁移进度。

迁移目标表结构可直接参考 0000_init.sql 中otel_2.otel_traces的 DDL,而app/opentelemetry-collector/目录(config/base.yamlclickhouse/default_ddl/)则是新链路(SDK 直发 OTLP)的接收端,两套数据最终汇聚到同一张表,供/v4/traces/v4/meterics端点(见 API README)统一查询。

小结

AgentOps Exporter 是一次"关系型监控数据 → OTEL 化"迁移的典型工程实现:

  • 语义映射层(models.py)把 sessions/agents/actions/llms/tools/errors 六张表精确映射为 Trace/Span 树,ID 复用 session_id 与事件主键,LLM 数据附加 GenAI 语义约定属性;
  • 实时双写层(export.py + v2.py)在遗留 API 每个写端点内同步影子写入 ClickHouse,保证增量数据零延迟进入新架构;
  • 批量迁移层(processor.py)提供分页、并行读(21 路)+ 并发写(42 任务池)、240 秒写超时、断点续跑与失败记录审计的完整批量工具,并用DRY_RUNMAX_PAGES支持先小样本(万行级)验证、后全量执行的稳妥节奏。

从源码结构看,agentops.project.id属性贯穿写入、物化列与索引设计,是这一迁移方案中多租户隔离的关键约定;理解了这条链路,就能完整复现 AgentOps 从 Postgres 到 ClickHouse OTEL 后端的数据搬迁全过程。

【免费下载链接】agentopsPython SDK for AI agent monitoring, LLM cost tracking, benchmarking, and more. Integrates with most LLMs and agent frameworks including CrewAI, Agno, OpenAI Agents SDK, Langchain, Autogen, AG2, and CamelAI项目地址: https://gitcode.com/GitHub_Trending/ag/agentops

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询