使用 dlt 将 Langfuse 可观测性数据导出到数据湖仓
【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy 🛠️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt
导读
Langfuse 是开源的 LLM 可观测性与评估平台,它把 traces、评估数据集、评分与成本数据全部持久化在自托管的 PostgreSQL 中。本指南讲解如何借助 dlt(data load tool)内置的sql_database()源,将 Langfuse 的 Postgres 后端一键导出到 DuckDB 等数据湖仓/数仓,从而支撑离线分析、报表与批量评估。读完本文,你将掌握连接凭据配置、resolve_foreign_keys外键关系自动关联、table_names选择性加载等核心能力,并能独立搭建一条 Langfuse → 目标存储的完整数据管道。
Langfuse 数据为什么要导出
Langfuse 是一个开源 LLM 可观测性与评估平台,它会捕获 AI 应用的 traces 与 observations,管理评估数据集,并跟踪 LLM 调用成本——所有这些数据都存放在一个可自托管的、基于 PostgreSQL 的存储中(自托管场景下)。它的 Web UI 适合在线查看与人工评估,但要支持以下场景,把数据搬出到数据湖仓/数仓是必经之路:
- 分析与报表:跨项目统计 token 消耗、成本趋势、模型分布;
- 离线评估:把
datasets/dataset_items与评分结果批量拉入分析引擎,做批量评测; - 数据治理与审计:
audit_logs、comments等元数据进入统一数仓,与业务数据打通。
Langfuse 的关系型元数据全部持久化在PostgreSQL中,而 dlt 内置的sql_database()源让抽取变得非常直接:
- 用连接字符串直连 Langfuse 的 Postgres 后端;
- 设置
resolve_foreign_keys=True,让 dlt 自动把子表(如datasets → projects)关联起来; - 加载到DuckDB实现零配置的本地分析,或者换成任意其他 dlt 目标。
凭据配置:连接 Langfuse 的 Postgres 后端
在项目目录下新建.dlt/secrets.toml,填入你部署实例对应的值:
# .dlt/secrets.toml [sources.sql_database.credentials] drivername = "postgresql" host = "localhost" port = 5432 database = "langfuse" username = "langfuse" password = "langfuse"这些字段与仓库中ConnectionStringCredentials的配置规范一一对应(drivername、host、port、database、username、password、query),password会被 dlt 作为密文处理。[sources.sql_database.credentials]这一节名与sql_database()默认参数credentials=dlt.secrets.value的解析路径一致,你也可以改用完整连接字符串(如postgresql://langfuse:langfuse@localhost:5432/langfuse)通过dlt.secrets注入。
另外还需要安装 psycopg2 驱动:
pip install psycopg2-binary各部署形态下这些值从哪里找
- Docker Compose:查看
docker-compose.yml中的db服务。凭据由POSTGRES_USER、POSTGRES_PASSWORD、POSTGRES_DB环境变量设定;host 是服务名(若端口已发布到宿主机则用localhost)。 - Helm chart:查看
values.yaml中postgresql.auth配置,或 Langfuse 部署上的DATABASE_URL环境变量。 - Langfuse Cloud:云托管无法直接访问数据库,请改用 Langfuse API 进行导出,而不是直连 Postgres。
核心原理:sql_database()是如何工作的
本文方案的核心是 dlt 内置的 SQL 数据库源sql_database(),其实现位于dlt/sources/sql_database/__init__.py。调用它时,dlt 会:
- 反射(Reflect)数据库结构:通过 SQLAlchemy 连接数据库,读取 schema 中全部表/视图的列、主键、外键约束,无需手写任何建表或映射代码(对应源码中的
metadata.reflect(...)调用,__init__.py); - 生成资源(Resources):为每个表生成一个独立的
DltResource,逐表执行SELECT并按chunk_size(默认 50000 行)分批产出数据; - 推断并固化 schema:把 SQLAlchemy 列类型映射为 dlt 目标 schema 类型(见下文“类型推断”),在抽取的同时完成类型转换。
关键参数:reflection_level 与 resolve_foreign_keys
Langfuse 示例中使用了两个值得深入理解的参数:
reflection_level="minimal":只反射表名、可空性和主键,列的数据类型从实际数据中推断。可选值见ReflectionLevel定义:"minimal"/"full"(默认,在 minimal 基础上反射类型并做必要强转)/"full_with_precision"(额外携带 decimal、text、binary 的精度与 scale)。Langfuse 的input、metadata等 JSON 字段在minimal下会被 dlt 规范化拆分为子表,保持目标 schema 的稳定与可读。resolve_foreign_keys=True:把同 schema 内的外键翻译成 dlt 的references表提示(hint),从而在目标端生成可查询的表关系。其底层实现是get_table_references():遍历每个表的foreign_key_constraints,解析被引用表名与列,合并指向同一张表的多个外键,最终形成{"referenced_table": ..., "referenced_columns": [...], "columns": [...]}结构。需要留意的是,该开关会触发额外数据库调用去反射所有被引用表,因此会略微增加启动开销。
类型推断
在非minimal反射级别下,sqla_col_to_column_schema()负责把 SQLAlchemy 类型映射到 dlt 类型:Numeric→decimal/double、Integer→bigint、String→text、DateTime→timestamp(并透传 timezone 提示)、JSON/ARRAY→json、Boolean→bool、UUID→text。即使使用minimal级别,目标端建表时 dlt 仍会以这些推断结果为准,保证 DuckDB 等目标的列类型贴合源库。
Langfuse 数据模型:自动发现与关键表
dlt 会自动发现 Langfuse 数据库中的每一张表,并从外键推断它们之间的关系,无需任何 schema 配置。实际得到的表集合取决于你使用了哪些 Langfuse 功能——只有被使用过的功能,其对应的表才会有数据被写入。
与 LLM 可观测性与评估最相关的表
organizations— 顶层多租户容器,每个 org 拥有自己的 projects 与 members。projects— 一个 project 划定某个应用的所有 traces、scores 与 datasets。trace_sessions— 把相关 traces 归组为一个 session(如多轮对话),带environment字段用于区分 staging 与 production。datasets— 用于评估运行的精选示例集合,可通过 Langfuse Web UI 或 SDK 创建。dataset_items— dataset 内的单条记录,每条携带input、可选的expected_output,以及回指到产生它的 source trace/observation 的链接。score_configs— 命名评分类型定义(数值或分类),可带min_value/max_value边界;子表score_configs__categories存放分类标签。eval_templates— 带版本的 LLM-as-judge 模板,包含prompt、model、provider和结构化输出 schema;子表eval_templates__vars存放变量名。models— LLM 模型定义,含input_price、output_price与 tokenizer 配置,用于计算每条 trace 的成本。prices/pricing_tiers— 关联到模型定义的分层定价规则。annotation_queues— 人工标注工作流;条目记录在annotation_queue_items,指派人记录在annotation_queue_assignments。comments— 用户附加到任意 Langfuse 对象(trace、dataset 等)上的评论。dashboards/dashboard_widgets— 已保存的分析仪表盘及其图表配置。audit_logs— 平台内所有 create/update/delete 动作的完整审计轨迹。
部分表属于 Langfuse 内部实现(如_prisma_migrations、background_migrations),可以通过sql_database(..., table_names=...)显式选择要加载的表子集,把它们排除在外。
目标 schema 实体关系图
dlt 在minimal反射级别下加载后,会在目标端(如 DuckDB)生成如下 ER 图所示的表结构(包含 dlt 自身的_dlt_loads、_dlt_version、_dlt_pipeline_state元数据表,以及由 JSON 列规范化拆分出的*__*子表与_dlt_parent_id/_dlt_list_idx关联列):
从上图可以直观看到:projects与organizations、datasets、score_configs、trace_sessions、eval_templates、comments等表之间存在外键关联——这正是resolve_foreign_keys=True发挥作用的地方;dataset_items的input、expected_output、metadata等 JSON 字段则被展开为dataset_items__input__messages、dataset_items__expected_output__parts等多层子表。
完整源码:Langfuse → DuckDB 导出管道
以下脚本是完整的可运行示例。它定义了一个langfuse_source源,把整库导出到本地 DuckDB:
import dlt from dlt.sources.sql_database import sql_database @dlt.source def langfuse_source(credentials=dlt.secrets.value): return sql_database( credentials=credentials, reflection_level="minimal", resolve_foreign_keys=True, ) if __name__ == "__main__": pipeline = dlt.pipeline( pipeline_name="langfuse", # can be configure to any dlt destination destination="duckdb", ) load_info = pipeline.run(langfuse_source) print(load_info)关键点解读:
@dlt.source装饰器把langfuse_source声明为 dlt 源;credentials=dlt.secrets.value让它从.dlt/secrets.toml自动读取[sources.sql_database.credentials]一节。reflection_level="minimal"只做轻量反射、类型由数据推断,配合 JSON 列自动拆分,适合 Langfuse 这种含大量嵌套input/metadata字段的库。resolve_foreign_keys=True让datasets → projects、prices → models等外键关系在目标端以references提示保留。destination="duckdb"零配置即可本地运行;注释已标明可以换成任意 dlt 目标(如postgres、filesystem、bigquery、snowflake等)。
运行后,load_info会打印每个资源(表)的加载状态与行数。管道默认把数据写入 DuckDB 的langfusedataset,dlt 会自动建好目标表与_dlt_loads等元数据表。
进阶配置与常见优化
1. 只加载关键表,跳过内部表
Langfuse 库中包含_prisma_migrations、background_migrations等内部表。通过table_names显式指定子集即可只加载与可观测性相关的表:
source = sql_database( credentials=dlt.secrets.value, reflection_level="minimal", resolve_foreign_keys=True, table_names=[ "organizations", "projects", "trace_sessions", "datasets", "dataset_items", "score_configs", "eval_templates", "models", "prices", "pricing_tiers", "comments", "audit_logs", ], )从源码看,sql_database()在table_names给定后只会反射并生成这些表的资源;而如果使用.with_resources(...),则会先反射全库再做过滤,对大 schema 而言显式传table_names能显著提升性能。
2. 按增量方式持续同步
Langfuse 各表普遍带有created_at/updated_at时间戳列。使用sql_table资源配合dlt.sources.incremental即可实现增量抽取,例如只加载增量更新的dataset_items:
import dlt from dlt.sources.sql_database import sql_table table = sql_table( table="dataset_items", incremental=dlt.sources.incremental("updated_at"), reflection_level="full", resolve_foreign_keys=True, )dlt 会在状态(state)中记录游标最大值,后续运行自动生成WHERE updated_at >= ...的查询,避免全表重复抽取。
3. 切换后端以提升吞吐
sql_database()支持四种后端(backend参数,见 源码定义):sqlalchemy(默认,产出字典批次,无额外依赖)、pyarrow(产出 Arrow 表,类型保真、性能更好)、pandas(产出 DataFrame)、connectorx(Rust 实现,通常最快但会忽略chunk_size)。对包含decimal价格的models/prices表,PyArrow 后端能避免精度损失:
source = sql_database( credentials=dlt.secrets.value, reflection_level="full_with_precision", resolve_foreign_keys=True, backend="pyarrow", )4. 在配置文件中管理表与列选择
除了在代码中传参,也可以在config.toml中声明要加载的表与列(表名列名大小写敏感,须与数据库中完全一致):
# 选择要加载的表 [sources.sql_database] table_names = [ "projects", "dataset_items", ] # 只加载 dataset_items 的部分列 [sources.sql_database.dataset_items] included_columns = [ "id", "dataset_id", "input", "expected_output", "status", "updated_at", ]5. 扩展到其他目标
DuckDB 适合本地零配置分析;生产环境中把pipeline的destination换成"postgres"、"bigquery"、"snowflake"、"filesystem"(parquet/iceberg/delta 湖)等即可,源端代码无需改动。若 Langfuse 部署在远程且数据库不可直连,可参考仓库中sql_database官方源文档 里的 SSH 隧道方案(sshtunnel+ 自定义 SQLAlchemy engine)来建立安全通道。
总结
Langfuse 的 Postgres 后端包含了 traces、scores、datasets、成本模型等完整的 LLM 可观测性数据,通过 dlt 内置的sql_database()源可以在零手工建表的前提下把这些数据整库或按需导入 DuckDB 及其他数据湖仓/数仓。核心要点可归纳为三条:在.dlt/secrets.toml配置postgresql凭据并安装psycopg2-binary;用reflection_level="minimal"获得稳定的嵌套 JSON 规范化、用resolve_foreign_keys=True保留表间关系;用table_names或incremental控制加载范围与增量同步。相关源码与测试位于仓库的dlt/sources/sql_database/与tests/sources/sql_database/,官方源的完整用法可进一步查阅sql_database使用与配置文档。
【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy 🛠️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考