Lightdash 内部使用分析架构:从事件流到 Parquet 再到系统 Explore 的完整技术解析
2026/9/18 8:34:04 网站建设 项目流程

Lightdash 内部使用分析架构:从事件流到 Parquet 再到系统 Explore 的完整技术解析

【免费下载链接】lightdashAgentic BI. Analytics at the speed of code ⚡️项目地址: https://gitcode.com/GitHub_Trending/li/lightdash

导读

本文是 Lightdash 开源仓库中 docs/usage-analytics/architecture.md 的深度技术解读。它讲解了 Lightdash 内部使用分析(Internal usage analytics)的完整数据链路:后端事件投影如何经缓冲写入与定时压缩变成可查询的 Parquet 文件,后端如何通过服务端持有的凭证签发短期 GET 签名 URL 让 DuckDB 读取,以及围绕组织隔离、凭证信任边界展开的安全设计。读完本文,你将掌握 Lightdash 使用分析项目的端到端架构、三个核心模块(连接器、凭证提供器、项目供给栈)的实现原理,以及该功能在生产上线前仍需补齐的边界。

注意:本文描述的是后端托管的内部元数据项目(与组织的业务模型完全隔离),当前处于功能开关(feature flag)控制的内部预览阶段,并非面向客户的正式交付方案。用户使用 Lightdash 常规的 Explore/查询流程探索固定维度和指标,它不是一组硬编码报表查询,且不依赖 dbt,也不需要 MotherDuck 账号

当前实现与上线意图的差距

架构文档在一开始就用对照表划清了"已实现"与"生产化仍需"的边界,这对评估任何内部功能都至关重要:

领域已实现生产化仍需
项目内部PREVIEW项目,带provisioning_source=analytics标记;包含 Query Events 与 AI Usage 两个数据源最终元数据项目的生命周期与角色设计
入口会话认证、仅组织管理员可用的 create-or-get 端点管理员导航/按钮,无连接器配置 UI
启用方式标准analytics-project特性解析器:Console 组织开关或部署 ENV 开关面向客户推广前的受控线上验证
源组织持久化、已授权的项目组织;忽略旧版本地覆盖共享实例隔离的线上验证
存储认证复用后端保留的写入凭证;签名 GET URL 交给 DuckDB验证实际 IAM 效果;只读加固延后(PROD-11103)
模型后端拥有的query_eventsai_usage两个 Explore重新评估其他事件流、元数据丰富与更多事件覆盖

文档明确警告:历史 ticket 和交接提案中可能描述过不同设计,它们不能作为生产供给、资源名丰富、导出或最终角色限制已实现的证据

端到端数据流总览

整个链路可以用文档中的流程图概括:

Backend event projections → buffered writer: gzip JSONL in events/raw/ (org / stream / date) → scheduled compactor: typed Parquet in events/compacted/ (same partitions) Signed-in org admin → create-or-get endpoint → fixed system explores stored in the app database → Explore query → backend authorization and org-scoped file discovery → exact-file signed GET URLs → isolated in-memory DuckDB → query results

前半段是"写入路径"(事件捕获与压缩),后半段是"读取路径"(项目供给与查询执行)。这两段在代码中分属不同的模块,下面分别深入。

写入路径:缓冲写入器与分区

事件注册表(Event Registry)

事件注册表 是整个使用分析流的"白名单":它限定哪些后端事件会被投影进使用事件流,并为每条流定义压缩后的 Parquet schema。核心结构如下:

  • ProjectedEvent:所有被投影事件的联合类型,目前包括QueryCompletedEventAiUsageEventAiAgentStepCompletedEventAiAgentToolCallCompletedEventDataAppStreamEvent
  • eventStreamRegistry:按事件名分派的投影函数表,由queryEventsProjectionsaiUsageProjectionsagentStepsProjectionsdataAppEventsProjections合并而成——不在白名单中的事件会被 sink 直接忽略
  • compactedStreamSchemas:每条流对应的类型化 Parquet 列定义(query_eventsai_usageagent_stepsdata_app_events),供夜间压缩任务使用;raw 区出现但未在此注册的流会被跳过(其 raw 文件永不删除)

从源码注释可以确认:新增一条流/事件 = 在事件类型中加入该事件并在下方补一个投影条目,其他部分无需改动。

缓冲写入器(Buffered Event Stream Writer)

BufferedEventStreamWriter 是每进程一个的写入器,负责把内存中的事件行批量刷入 S3。它的关键行为:

  • (org_id, stream, dt)三元组把缓冲中的行分组(dt由事件时间戳取 UTC 日期,见getUtcDate);
  • 每组写成一个独立的 gzip JSONL 对象,键格式为events/raw/org_id=<org>/stream=<stream>/dt=<YYYY-MM-DD>/<writerId>-<uuid>.jsonl.gz,其中writerId是进程启动时生成的 8 位随机 UUID;
  • 采用"批量阈值 + 定时器"双触发:flushBatchSize达到即刷,flushIntervalMs定时兜底;缓冲满或写入器已关闭时丢弃事件并递增UsageEventsDropped指标;
  • 写入器从不向共享 Parquet 对象追加数据,每次都是独立 PUT;
  • PUT 失败最多重试 2 次,仍失败则整组丢弃并告警;flush()永不向调用方抛异常(内部 catch 后仅记录 warn 日志)。

该实现从源头说明了"小文件问题"的成因:高并发下每进程独立写 JSONL,会产生大量细碎对象,这正是引入压缩器的主要原因。

压缩器(Usage Events Compactor)

UsageEventsCompactor 负责把已闭合的分区转换为列式 Parquet。关键设计点:

  • 只处理闭合分区dt严格早于今天 UTC 的分区才参与压缩(groupRawKeysIntoPartitionsparsed.dt >= todayUtc直接跳过),今天的 raw 事件不进入读取路径;
  • 确定性 part 文件名buildPartFileName对排序后的 raw key 列表做 SHA-1 哈希生成part-<hash>.parquet,因此对同一输入集重跑会覆盖同一 part 而非产生重复数据;迟到文件会改变 key 集合,从而产生额外的 part——所以并非一天一个文件、一个组织一个文件;
  • 精确读取,绝不 globbuildCompactionSql把列出的原始文件逐个拼进read_json的显式文件列表(而非通配符),避免"列list之后写入的对象"被未压缩删除或重复压缩;
  • 转换成功后才删除 raw 对象:删除采用批量(每批 1000 个)+ 最多 3 次重试 + 指数退避;若删除最终失败会显式抛错并给出人工清理指引——因为下次运行会因 key 集合不同生成第二个 part,造成行重复;
  • 未知流跳过:schema 未注册的分区只计数并告警,raw 文件原样保留;
  • 运行上限:单次最多处理 500 个闭合分区(MAX_PARTITIONS_PER_RUN),超出的留待下轮,通过UsageEventsCompactionBacklog增长型指标反映落后程度。

压缩把独立写入的小文件开销降了下来,产出适合分析的列式格式。但读取端只认压缩输出,因此捕获事件不会立刻出现在 Explore 中;同时该栈不改变写入器或夜间流程,也不提供 exactly-once 投递保证。

项目创建与系统模型

create-or-get 端点

ProjectService.ensureAnalyticsProject支撑POST /api/v1/org/analytics-project,执行五步:

  1. 鉴权:校验会话的组织、组织管理权限与特性开关;请求体中不接受任何 org、bucket、路径或凭证参数;
  2. 先验证签名 Parquet 可读再建项目:通过createAnalyticsClient创建客户端并调用analyticsClient.test()做一次存储往返。模型编译在内存中完成,但该端点包含这次存储往返以确保不会为坏配置建出死项目;
  3. 获取 per-org 咨询锁runInAnalyticsProvisioningLock,数据库级 advisory lock),并查找provisioningSource === 'analytics'的既有项目;
  4. 复用或新建:不存在则用createWithoutCompile创建内部 DuckDB analytics 连接(type: DUCKDBconnectionType: ANALYTICSdatabase: 'memory'schema: 'main')的PREVIEW项目,然后saveExploresToCache保存两个编译好的系统 Explore;重复调用会刷新模型而不会删除已保存内容,还能修复上次模型保存失败的状态
  5. 返回{ projectUuid, url, created },前端重定向到返回的 URL。

身份识别来自内部标记 + 组织,而非显示名或 slug。slug 由 generateUniqueProjectSlug 分配:由 "Lightdash analytics" 派生出lightdash-analytics,被占用则依次尝试lightdash-analytics-1-2……重用的项目保留原 slug 与 UUID;冲突的普通项目绝不会被改造成 analytics 项目

系统 Explore 的编译

createAnalyticsExplorescompactedStreamSchemas派生维度、从 systemStreamMetrics 派生指标,明确只包含query_eventsai_usage(当前代码还包含data_app_events,但文档所述的供给列表以query_eventsai_usage为准——注册另一个写入器流不会自动暴露新的 Explore)。

关键机制:

  • 列类型映射:VARCHAR→STRINGTIMESTAMP→TIMESTAMPBOOLEAN→BOOLEANINTEGER/BIGINT→NUMBER
  • 编译后的模型sqlTable指向表名而非签名 URLdatabase: 'memory'schema: 'main'),用户 SQL 与聚合完全由用户在固定模型上的 Explore 选择生成;
  • 预置指标示例(query_events流):total_queries(COUNT query_id)、unique_users(COUNT_DISTINCT user_id)、avg_warehouse_execution_time_ms(AVERAGE)、p90_warehouse_execution_time_ms(90 百分位)等;ai_usagedata_app_events流同样有对应的计数、去重与过滤指标定义。

凭证与信任边界:三种截然不同的访问形式

文档用一张表强调:有三层访问,绝不可混淆

凭证/权威范围
用户 → Lightdash认证会话 + 组织授权 + 特性门访问该组织的 analytics 项目
后端 → 对象存储服务端持有的 usage-events 访问 key/secret(当前复用于摄取)潜在较宽的桶访问,含写权限;不按组织隔离
DuckDB → Parquet后端生成、短期的签名 GET URL精确选中的对象、HTTP 方法与有效期

当前验证过的云路径走 GCS 的 S3 兼容接口 + HMAC 凭证 + AWS S3 SDK。不运行云 CLI、不获取开发者 OAuth token、不供给身份、不需要 MotherDuck token

配置与生命周期

parseUsageEventsS3Config读取:

  • USAGE_EVENTS_S3_ENDPOINT(回退到基础S3_ENDPOINT
  • USAGE_EVENTS_S3_BUCKET(回退S3_BUCKET
  • USAGE_EVENTS_S3_REGION(回退S3_REGION
  • USAGE_EVENTS_S3_ACCESS_KEY/USAGE_EVENTS_S3_SECRET_KEY(回退到S3_ACCESS_KEY/S3_SECRET_KEY

基础 S3 配置必须有效parseBaseS3Config返回 null 时整个配置为 null)。这些是服务器配置而非用户可编辑的项目凭证。文档特别指出:当 analytics 用 GCS、本地普通存储用 MinIO 时,用 endpoint 覆盖来避免重定向普通存储;该覆盖同样作用于 usage-events 写入器/压缩器配置,而不仅是读取端

源解析器(S3 Analytics Source)

S3AnalyticsSource是读取路径的核心:

  • 只列events/compacted/org_id=<validated-org>/前缀,并校验每个返回 key 确实以该前缀开头(越界立即抛错);
  • 用正则^stream=(query_events|ai_usage|data_app_events)\/dt=(\d{4}-\d{2}-\d{2})\/[a-zA-Z0-9_-]+\.parquet$只挑选受支持流的 Parquet 路径,覆盖所有保留日期;日期过滤属于 Explore 查询,不是固定的源窗口——不会列整个桶再事后按 org 过滤;
  • 分页有界:每页最多 1000 个对象、最多 100 页;文件总数上限 10,000(MAX_FILES),超限即抛错而非静默截断;畸形/不完整列表或空 manifest(无任何表)均 fail closed;
  • 为每个解析结果精确签名GetObject请求,有效期 900 秒(15 分钟),并按表名分组排序返回;
  • 每次解析都新建 SDK 客户端、结束时销毁finally { client.destroy() }),降低凭证驻留面——但文档提醒:销毁客户端不等于擦除进程配置或吊销凭证,后端仍持有配置中的凭证;
  • endpoint 安全校验:远程必须 HTTPS;HTTP 仅允许 loopback(localhost / 127.0.0.1 / [::1]);URL 不允许带用户名、密码、query 或 fragment,path 必须是/
  • 出错时统一抛通用错误消息,避免 SDK 错误泄露凭证、签名或对象内容。

每个新的 DuckDB 会话都会重新解析 manifest 与签名。URL 过期即失败;新查询/新会话会自动获得新签名。凭证轮换必须通过部署的正常 secret/重启机制刷新后端配置,这不是自动轮换实现。

什么能进入 DuckDB

签名 manifest 的内容边界

发给 DuckDB 的签名 manifest 包含:表名、预期前缀(scope)、精确 URL 列表。它不包含桶访问 key/secret,也不包含任何 bearer token。但签名 URL 本身就是一种 bearer 能力:任何持有者在有效期内都能读该对象。文档明确要求:绝不可把签名 URL 放进日志、API 响应、ticket 或持久化模型中

DuckDB 内部读取器(DuckdbWarehouseClient)

DuckdbWarehouseClient配合duckdb_parquet连接类型(见 analyticsProjectClient)实现隔离读取:

  • 每次查询创建隔离的内存实例(256 MB DuckDB 内存上限、32 线程以重叠远端 footer 与列读取);显式调用方资源限制会覆盖这些内部读取默认值,其余 DuckDB 默认值不变;
  • 远程存储强制 HTTPS;HTTP 仅限 loopback 测试端点;
  • 校验规范路径在受信 scope 内;禁止任意 globbing/路径替换,禁止把签名 URL 与宽范围存储 secret 混用
  • 设置精确allowed_paths禁用通用外部访问与磁盘溢出(spill)
  • 仅在私有查询实例内部启用 HTTP 元数据、Parquet 元数据与外部文件缓存;实例成功或失败都会关闭,缓存字节与签名 URL 都不会被其他请求或组织复用
  • 仅针对这些精确对象构建read_parquet临时视图;
  • 限制用户 SQL,阻断可能泄露视图 SQL 的 catalog 访问;禁用 profiling、清理原生查询错误。

schema 绑定保留union_by_name=true,因此旧文件可以缺少新列。缓存避免了绑定、校验、执行期间重复读取远端元数据。但这不改变"全历史发现"、不跳过旧文件、也不放宽 10,000 文件上限;大历史仍需扩展工作(PROD-11111),这些设置不保证任意数据量下的延迟上界。每个并发查询有各自线程预算;部署级与 per-org 并发限制在客户推广前仍需验证

需要澄清:查询级缓存生命周期不意味着 Lightdash 没有持久化查询结果。常规 result/history 路径依然存在且需授权;analytics 检查覆盖项目访问与 result/history 检索;本预览中导出与定时下载被阻断。同时,关掉 flag 是拒绝而非删除,也不是对已签发能力或已返回数据的加密撤销。

组织隔离:已实现的边界与剩余风险

analyticsProjectClient在服务层授权检查之后,只为持久化项目的 org读取文件;所有环境中都忽略旧版本地 org/source-org 覆盖。生产执行使用面向目标组织的标准特性解析器(isAnalyticsProjectEnabledorganizationUuid读取FeatureFlags.AnalyticsProject)。

两个必须讲清的边界:

  1. 前缀过滤是后端授权边界,不是桶 IAM。存储签名把"精确对象 + HTTP 方法"能力下发给 DuckDB,但受信的签名者仍可用更宽的 key 访问或签名其他对象。签名者/后端被攻破不在这些 URL 的保护范围之内;路径篡改测试不能证明任意 SQL 或后端攻破是安全的。
  2. 只读凭证降低签名者的写/删权限暴露,但本身不能阻止跨 org 前缀的读。延后的身份设计是"独立部署级读取身份 + 复核的应用侧 org 绑定",而非按应用组织建服务账号;在写入者身份上再建一把 key 并不会收窄权限。PROD-11103 记录了基础设施指引,但本栈不含任何 Terraform/IAM 变更或只读切换

运维、验证与上线准备

本地与线上测试

  • 本地测试:见 local-testing.md。在隔离开发实例的.env.development.local中设置LIGHTDASH_ENABLE_FEATURE_FLAGS=analytics-project(显式禁用优先);以管理员身份打开Organization settings → Lightdash analytics (Beta)创建/打开项目,或用浏览器控制台直接调用POST /api/v1/org/analytics-project(端点不接受任何配置,身份来自会话)。相关端点全部位于AnalyticsProjectController,由AnalyticsProjectService支撑,均要求会话认证的 org 管理员 + analytics 特性。
  • 凭证验证:见 credentials.md。MinIO 环回测试命令为:
    ANALYTICS_S3_SMOKE_ENDPOINT=http://localhost:9000 pnpm -F backend test src/services/ProjectService/analyticsProject/S3AnalyticsSource.smoke.test.ts

    线上只读测试(需ANALYTICS_S3_LIVE_ENV_FILE指向含USAGE_EVENTS_S3_*的安全环境文件,以及ANALYTICS_S3_LIVE_ORG_UUID):

    pnpm -F backend test src/services/ProjectService/analyticsProject/S3AnalyticsSource.live.test.ts

    两个外部测试套件未显式配置时都会被跳过;绝不提交环境文件、签名 URL、凭证或客户专属测试夹具。单测覆盖前缀过滤、分页、精确 GET 签名、无效 scope、无凭证 manifest、原生错误脱敏、被阻断的 catalog/文件读取。

  • 线上启用:见 live-testing.md。NODE_ENV=production下可用既有USAGE_EVENTS_S3_*配置运行;在 Console 为目标组织启用analytics-project(无缓存、免重启),或加入LIGHTDASH_ENABLE_FEATURE_FLAGS全局启用(需部署/重启)。标准优先级:ENV 启用 > ENV 禁用 > 数据库组织覆盖/默认。回滚 = 关掉 Console 组织开关或移除 ENV 启用并显式禁用;已签发 URL 最长 15 分钟内仍有效,flag 关闭不会撤销已返回的数据。

排障清单

  • 缺数据:依次检查捕获(capture)、压缩是否成功、源组织、Explore 日期过滤、保留期与流 schema。无受支持文件会 fail closed;部分缺失的流与 schema 演化仍需上线设计。
  • 访问错误:检查端点、桶与已配置 key 的权限——不要打印 key、SDK 请求细节或签名 URL。新查询可从 URL 过期恢复;它无法修复已吊销或错误 scope 的源凭证。
  • 供给延迟:包含存储访问。大组织历史可能撞上列文件上限;上限是"失败"而非"静默只暴露部分历史"。历史可用性受保留文件限制而非文件年龄;一年期负载仍未基准化(PROD-11111)。
  • 上线前必查:生产授权、角色限制、有效 IAM、查询/缓存结果面;专用只读凭证是可选的延后加固(PROD-11103);隐藏 UI 与文件夹隔离不能替代访问控制

文档如实说明:读取路径已在真实 GCS 上以原生 DuckDB 与本地 API 对两个 Explore 做过验证,本地 UI 供给经人工确认;聚焦测试覆盖复用、访问拒绝与凭证边界。但它不是生产多租户安全审计,也不声称每个仓库提供方与故障模式都已验证

相关工单与演进线索

  • PROD-11059:实现/分诊(本地 triage 主 ticket);
  • PROD-8603:管道架构;
  • PROD-11103:延后的专用只读凭证;
  • PROD-11111:大历史与并发查询扩展;
  • PROD-11152:内容版本检测与更广的内容同步;
  • PR 栈演进:连接器 + 特性开关 → 签名 URL 桶访问 → 项目供给端点(含两个系统模型)→ 本文档所属的内部文档。

总结

Lightdash 内部使用分析是一套"后端拥有元数据项目"的工程栈:事件注册表白名单 → 缓冲 JSONL 写入 → 定时压缩为分区 Parquet → 按 org 前缀发现 + 精确签名 GET URL → 隔离内存 DuckDB 执行。它的价值在于把"平台自身的使用行为"变成用 Lightdash 常规 Explore 流程即可分析的固定模型,同时明确划出了三条信任边界(用户-会话、后端-存储、DuckDB-签名 URL)。任何在生产中复现该方案的人都应把 architecture.md、local-testing.md、credentials.md 与 live-testing.md 四份文档作为配套阅读,并严格遵守其中"未实现即未实现"的边界声明。

【免费下载链接】lightdashAgentic BI. Analytics at the speed of code ⚡️项目地址: https://gitcode.com/GitHub_Trending/li/lightdash

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

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

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

立即咨询