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

📅 发布时间:2026/9/17 22:19:12
AgentOps Exporter:从 Supabase (Postgres) 到 ClickHouse 的 OpenTelemetry 数据迁移与双写通道
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/agentopsAgentOps 的后端正处于从传统关系型存储Supabase/Postgres向 OpenTelemetry ClickHouse 架构迁移的过程中app/api/agentops/exporter/目录就是这次数据搬迁的核心通道。本文基于该目录的 README 及其三个核心源码文件export.py、processor.py、models.py完整讲解 Exporter 的两类使用方式作为遗留 API 的实时双写处理器、作为可断点续跑的批量迁移工具以及底层 Session/Agent/Event → Trace/Span 的 OTEL 语义映射规则帮助读者理解 AgentOps 如何将 92.5 GB 级别的存量 Agent 会话数据转换为标准的 OTEL Traces 格式。Exporter 的定位与双重角色根据 Exporter README 的原始描述该服务承担两个职责一次性迁移One-Time Migration一个可按需执行的脚本把 Postgres 存量 schema 中的所有数据迁入 ClickHouse并行处理器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 → trace、events → span、agent → 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/sessions与POST /v2/create_sessionawait supabase.table(sessions).insert(session).execute()之后紧跟await export.create_session(session)见 v2.pyPOST /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.pyPOST /v2/update_events对每类事件执行 Supabase update 后调用对应的export.update_*函数。也就是说Supabase 写入是主路径ClickHouse 写入是并行影子流两条流共享同一份经过event_handlers处理后的数据结构。映射模型Session 变 Trace事件变 Spanmodels.py 用 Pydantic 定义了源端六张表到 OTEL 对象的映射。EXPORT_TABLES_MODELS见 processor.py给出了完整的表-模型对照Supabase 表Pydantic 模型OTEL 目标形态SpanKindsessionsSessionTrace 的根 Spanspan_id trace_id session.idSERVERagentsAgent父级 Spanparent_span_id session_idINTERNALactionsActionEvent子 Spanparent_span_id agent_idINTERNALllmsLLMEvent子 Spanparent_span_id agent_idCLIENTtoolsToolEvent子 Spanparent_span_id agent_idINTERNALerrorsErrorEvent点事件 Spanparent_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_idspan 的SpanId就是各表的主键 UUID无需生成新 ID天然满足幂等对齐时间单位转换ClickHouse 要求纳秒/DateTime64datetime_to_ns()把datetime转为纳秒把整型时间戳视为毫秒再乘 1e6见 models.pyLLM 事件采用 GenAI 语义约定LLMEvent.to_span()生成gen_ai.prompt.{i}.content、gen_ai.completion.{i}.finish_reason、gen_ai.usage.prompt_tokens等属性同时保留llm.model、llm.cost、llm.total_tokens等 AgentOps 私有属性并标注了 GenAI 语义约定的来源OpenTelemetry 官方规范见 models.pySpan 到 ClickHouse 行的序列化Span.to_clickhouse_dict()输出的列名与DESCRIBE otel_traces完全对齐Timestamp、TraceId、SpanId、ParentSpanId、SpanName、SpanKind、ResourceAttributes、Events.Timestamp等并在写入前强制注入ResourceAttributes[agentops.project.id] project_id见 models.py。这个agentops.project.id注入不是普通属性ClickHouse 侧的otel_traces表通过MATERIALIZED列直接派生project_id见 0000_init.sqlproject_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_id、idx_span_id、idx_project_id。因此 Exporter 正确写入agentops.project.id是多租户检索前端/v4/traces?project_id...能正常工作的前提。实时写入路径export.py 的调用链export.py 对外暴露一组异步函数供 API 层同步调用。以创建 LLM 事件为例调用链是校验data中必须携带session_id与agent_id_get_project_id(session_id)通过 Supabase 客户端反查sessions表拿到project_id构造 Pydantic 模型如LLMEvent(**data)调用to_span(trace_idsession_id, parent_span_idagent_id, project_idproject_id)生成Spanclickhouse_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_driver的escape_param见 processor.pyDRY_RUN True时只打印不写库便于灰度验证。批量迁移路径processor.py 的分页、并行与断点processor.py既是被export.py引用的内部库也是可独立运行的 CLI 入口if __name__ __main__时load_dotenv(.env, overrideTrue)后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.pyinit_files()初始化两个持久化文件dropped_records.csv记录迁移中因异常被丢弃的记录与last_session_id.txt断点检查点count_session_rows()统计剩余 sessions 行数total_pages row_count // BATCH_ROW_COUNT逐页调用process_page(offset, limit)get_sessions按id last_session_id ORDER BY id ASC LIMIT/OFFSET分页拉取UUID 按字典序排再以PARALLEL_READS21为批并行执行get_session_as_traceprocess_page维护一个最多MAX_CONCURRENT42的待写入任务池每条 trace 通过write_trace_with_timeoutasyncio.wait_for(..., timeout240)写入超时或异常的记录落到dropped_records.csv而不是中断整体迁移get_session_as_trace中每成功处理一个 session 就调用write_last_session_id把当前session_id回写到last_session_id.txtseek truncate write因此任务中断后重跑可从断点继续不会重复导出已完成的 sessionMAX_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_rows从otel_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 记录了一条完整的低风险验证路径结合源码可以落地为启动一个 dev 环境的 Supabase Postgres 实例并准备一个一次性的throwawayClickHouse 实例从生产库选取一段样本思路为SELECT ... LIMIT 0, 10000对应批量参数BATCH_ROW_COUNT1000MAX_PAGES限制页数写入 dev Supabase启动 exporter 处理实例python processor.py直接运行通过DRY_RUN先干跑确认输出结构再正式写入用dropped_records.csv核对失败记录用last_session_id.txt确认迁移进度。迁移目标表结构可直接参考 0000_init.sql 中otel_2.otel_traces的 DDL而app/opentelemetry-collector/目录config/base.yaml、clickhouse/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_RUN与MAX_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),仅供参考