DataHub Starburst Trino Usage 连接器:基于 Event Logger 摄取 Trino 使用统计的完整指南
2026/9/19 22:03:32 网站建设 项目流程

DataHub Starburst Trino Usage 连接器:基于 Event Logger 摄取 Trino 使用统计的完整指南

【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

导读

starburst-trino-usage是 DataHub 元数据摄取框架中面向生产环境的专用连接器,负责将 Starburst Trino 集群的查询使用情况(usage statistics)摄取进 DataHub,为数据集的热度分析、Top N 查询展示和用户行为洞察提供数据基础。本文将围绕该连接器的配置、能力边界与底层实现展开,读完你可以掌握如何编写可运行的摄取配方(recipe)、理解其数据采集原理,并能结合源码定位与排查常见问题。

模块概述:从 Trino 到 DataHub 的使用统计管线

在 DataHub 的摄取体系中,starburst-trino-usage是一个专注于**使用统计(Usage Stats)**的 source 模块。它不同于普通的元数据(Schema/Table)摄取,而是读取 Trino 集群内部记录的查询历史,将其加工为"哪个用户在什么时间、访问了哪些表、执行了什么 SQL"的结构化事件,最终以 DataHub 的 usage 工作单元(MetadataWorkUnit)写入目标端(如datahub-rest)。

从源码注释与实现看(starburst_trino_usage.py),该模块的核心数据来源是Starburst Event Logger 的审计日志表completed_queries,通过 SQLAlchemy 连接审计数据库(通常为 PostgreSQL),解析其中的accessed_metadataJSON 字段来还原表与列的访问明细,再按时间桶聚合后产出 usage 统计。

模块在连接器注册表中的登记信息(datahub.json)也印证了它的定位:

属性
连接器类型starburst-trino-usage
实现类datahub.ingestion.source.usage.starburst_trino_usage.TrinoUsageSource
平台 IDtrino
支持状态GA(生产可用)
能力(Capability)USAGE_STATS,默认启用

前提条件(Prerequisites)

官方文档(starburst-trino-usage_pre.md)明确列出运行摄取前必须满足的三个前提:

  1. 网络连通性:摄取节点必须能够访问 Trino 集群(或托管审计库的数据库)的网络端点,即配方中的host_port可达;
  2. 有效认证凭据:提供具有查询审计表权限的用户名与密码;
  3. 元数据 API 读权限:账号需要对 Event Logger 暴露的审计元数据(completed_queries表及相关字段)具备只读权限。

此外,由于实现依赖 Starburst Event Logger,运行环境需要满足:

  • Trino 侧已开启并配置好Event Logger,能够持续写入completed_queries审计表;
  • 审计表位于某个可被该连接器直接查询的 catalog/schema 下,通过配方中的audit_catalogaudit_schema指定;
  • 摄取进程所在环境已安装 DataHub ingestion 依赖(该模块代码位于datahub.ingestion.source.usage包内,随metadata-ingestion一起分发)。

快速开始:完整的摄取配方与参数逐项解析

仓库随文档附带了可直接参考的最小配方 starburst-trino-usage_recipe.yml,完整内容如下:

source: type: starburst-trino-usage config: # Coordinates host_port: yourtrinohost:port # The name of the catalog from getting the usage database: hive # Credentials username: trino_username password: trino_password email_domain: test.com audit_catalog: audit audit_schema: audit_schema sink: type: "datahub-rest" config: server: "http://localhost:8080"

其中各配置项的职责如下:

配置项必填说明
host_portTrino 服务地址(host:port),用于建立 SQLAlchemy 连接
database要采集使用统计的 catalog 名称。源码中用它过滤accessed_metadata,只保留与所选 catalog 一致的访问记录(见下文原理部分)
username/password访问审计数据库(Event Logger 后端)的认证凭据
email_domain当用户名不是完整邮箱时,追加到用户名后的邮箱域名,用于 DataHub 用户统计 UI 正确关联用户
audit_catalog存放审计表(completed_queries)的 catalog 名称
audit_schema存放审计表(completed_queries)的 schema 名称

这些字段与源码中的配置类TrinoUsageConfig(继承自TrinoConfigBaseUsageConfig)一一对应(starburst_trino_usage.py):

  • email_domain:将被追加到用户名的邮箱域名;
  • audit_catalog/audit_schema:审计表所在位置;
  • database:从中获取使用统计的 catalog 名。

继承自基类的进阶参数

除了配方中显式列出的字段,TrinoUsageConfig还继承了TrinoConfigBaseUsageConfig的大量参数,可在需要时按需追加,包括但不限于:

  • 时间窗口start_time/end_time,控制查询completed_queries的时间范围(默认值遵循BaseUsageConfig,源码在拼装 SQL 时通过strftime("%Y-%m-%d %H:%M:%S.%f %Z")格式化后注入查询);
  • 聚合粒度bucket_duration,决定按多长时间窗口聚合访问事件(源码通过get_time_bucketstarttime取桶);
  • 用户过滤user_email_pattern,在写入读事件时用于匹配用户邮箱;
  • 查询采样top_n_queriesinclude_top_n_queriesformat_sql_queriesqueries_character_limit,控制每条数据集产出 Top N 查询的数量与 SQL 格式化;
  • 平台实例与环境platform_instance/env,用于生成带实例限定符的 dataset URN;
  • 连接选项options,透传给create_engine的 SQLAlchemy 连接参数;
  • 表/视图开关include_tables/include_views,集成测试中同样验证了这些字段可被正确解析(见 test_starburst_trino_usage.py)。

运行方式

与 DataHub 其他连接器一致,保存好上述配方后,通过 CLI 执行:

datahub ingest -c starburst-trino-usage_recipe.yml

摄取完成后,可在 DataHub 前端的数据集页查看使用统计、Top N 查询与访问用户信息。

能力与限制(Capabilities & Limitations)

支持的能力

依据官方文档(starburst-trino-usage_post.md)与连接器注册表,该模块当前支持的能力为:

  • USAGE_STATS(使用统计):默认启用,无需额外配置即可获得查询使用情况。这是本模块唯一登记的能力,supported: true

"Important Capabilities" 能力表是判断某项功能是否支持、是否需要额外配置的权威来源(即上文 datahub.json 所登记的能力清单)。建议在接入前先对照该表确认需求是否被覆盖。

已知限制

模块行为受限于源平台的 API、权限与元数据暴露情况,结合源码可梳理出如下约束:

  1. 仅统计成功的 SELECT 查询:查询 SQL 硬编码了query_type = 'SELECT'query_state = 'FINISHED'的条件,DDL、失败查询、进行中查询都不会被计入;
  2. 依赖 Starburst Event Logger:必须存在填充了completed_queries表的审计 catalog/schema,非 Starburst 发行版或未开启 Event Logger 的 Trino 无法使用本模块;
  3. 仅聚合单一 catalog:通过database配置项过滤,accessed_metadata中 catalog 不等于该值的访问会被直接跳过;
  4. 系统查询被忽略:catalog 以$system@开头的查询会从统计中排除(源码第 250-257 行);
  5. 用户身份需要可解析为邮箱:源码中对无邮箱的用户会拼接{user}@{email_domain},无法识别用户名的场景会回退为unknown@{email_domain}
  6. accessed_metadata为空或create_time/usr缺失的事件会被跳过,并计入报告中的num_joined_access_events_skipped计数。

底层原理:源码级解析数据采集链路

理解了配置之后,深入 starburst_trino_usage.py 有助于正确判断部署形态与排查问题。整体处理链路为:

get_workunits_internal() ├─ _get_trino_history() # 1. 查询审计库 completed_queries 表 ├─ _get_joined_access_event() # 2. 解析 JSON、清洗并结构化事件 ├─ _aggregate_access_events() # 3. 按时间桶 + 数据集聚合 └─ _make_usage_stat() # 4. 产出 usage MetadataWorkUnit

1. 审计查询:读取 completed_queries

模块内置的查询模板trino_usage_sql_comment(源码第 39-55 行)针对 Starburst Event Logger 的 completed queries 审计表设计,核心逻辑为:

SELECT DISTINCT usr, query, "catalog", "schema", query_type, accessed_metadata, create_time, end_time FROM {audit_catalog}.{audit_schema}.completed_queries WHERE 1 = 1 AND query_type = 'SELECT' AND create_time >= timestamp '{start_time}' AND end_time < timestamp '{end_time}' AND query_state = 'FINISHED' ORDER BY end_time desc
  • {audit_catalog}/{audit_schema}来自配置,指向 Event Logger 审计表;
  • {start_time}/{end_time}来自BaseUsageConfig的时间窗口配置;
  • 查询结果若为空,源码会记录SQL Result is empty日志并直接终止摄取,避免空结果继续走后续流程。

2. 事件解析:accessed_metadata 的反序列化

审计表中的accessed_metadata是 JSON 文本,源码定义了TrinoAccessedMetadata模型(字段含catalogNameschematablecolumnsconnectorInfo)来承载它。_get_joined_access_event会:

  • 解析并校验create_time(缺失则跳过并计数);
  • accessed_metadata通过json.loads反序列化为结构化对象;
  • 校验用户字段usr(缺失则跳过);
  • 构造TrinoJoinedAccessEvent事件对象,任何解析异常都会计入num_joined_access_events_skipped

3. 聚合:时间桶 + 数据集维度

_aggregate_access_events是统计语义的核心:

  • 使用get_time_bucket(event.starttime, self.config.bucket_duration)将事件归入时间桶;
  • catalog.schema.table三元组作为数据集唯一标识(resource);
  • 过滤$system@系统 catalog 与非目标 catalog 的访问;
  • 对用户名做邮箱归一化:若usr本身是合法邮箱则原样保留,否则拼接email_domain
  • 调用agg_bucket.add_read_entry(username, query, columns, user_email_pattern=...)将"谁读、读什么、读了几列"写入聚合桶。

4. 产出:usage 工作单元

_make_usage_stat将聚合结果转换为MetadataWorkUnit,其中数据集 URN 通过make_dataset_urn_with_platform_instance生成,平台固定为trino,并应用platform_instanceenv;Top N 查询相关的top_n_queriesformat_sql_queriesinclude_top_n_queriesqueries_character_limit等参数在此处生效。

5. 报告与可观测性

TrinoUsageReportSourceReport基础上增加了num_joined_access_events_skipped计数。摄取出错或部分事件被跳过时,可通过该计数与日志(如Field accessed_metadata is empty. Skipping ....The username parameter is missing. Skipping ....)快速定位是数据缺失还是解析失败。

测试验证:如何确认模块行为

仓库提供了针对该模块的集成测试 test_starburst_trino_usage.py,包含两类用例:

  1. 配置解析测试test_trino_usage_config):验证TrinoUsageConfig能正确解析host_portdatabaseusernamepasswordemail_domainaudit_catalogaudit_schemainclude_viewsinclude_tables等字段,是校验配方字段书写的参考范本;
  2. 摄取端到端测试test_trino_usage_source):通过patch模拟_get_trino_history返回预置的访问事件,驱动完整Pipeline(source 为starburst-trino-usage,固定时间冻结在2021-08-24 09:00:00),再结合 MCE 快照辅助断言输出结果。

测试资源位于tests/integration/starburst-trino-usage/目录,可作为理解accessed_metadata事件结构、校验预期输出格式的参考样例。

故障排查(Troubleshooting)

官方文档给出的排查顺序(starburst-trino-usage_post.md)是:

  1. 先验证基础四要素:凭据(Credentials)是否有效、权限(Permissions)是否覆盖审计表、网络连通性(Connectivity)、范围过滤(Scope filters,即database/ 时间窗口)是否正确;
  2. 再检查摄取日志:针对 source 特有错误信息定位并调整配置。

结合源码,常见的具体问题与对策:

现象可能原因处置建议
日志出现SQL Result is empty,无任何输出时间窗口内无 FINISHED 的 SELECT 查询,或审计表未写入扩大start_time/end_time窗口;确认 Event Logger 已启用
统计中缺少某 catalog 的数据database未设为该 catalog,或访问记录中catalogName与配置不一致核对database与审计记录中的 catalog 名
用户显示为xxx@test.comunknown@...用户名非邮箱格式确认email_domain配置正确
事件被大量跳过(num_joined_access_events_skipped增长)create_time/usr缺失,或accessed_metadata为空/解析失败检查审计表数据完整性,查看跳过日志定位具体字段
连接失败host_port不可达或凭据错误先手工用客户端连接host_port验证连通性与权限

小结

starburst-trino-usage以 Starburst Event Logger 的completed_queries审计表为数据源,通过"查询审计库 → 解析 JSON → 时间桶聚合 → 生成 usage 工作单元"四步链路,为 DataHub 提供了开箱即用(GA 状态、USAGE_STATS 默认启用)的 Trino 使用统计摄取能力。接入时只需确保网络、凭据与审计表三个前提,编写包含host_portdatabaseaudit_catalogaudit_schemaemail_domain的配方即可运行;若需更精细的控制(时间窗口、Top N 查询、用户过滤等),可在继承自TrinoConfigBaseUsageConfig的参数中按需扩展。

延伸阅读

  • 模块官方文档:starburst-trino-usage_pre.md、starburst-trino-usage_post.md
  • 完整配方示例:starburst-trino-usage_recipe.yml
  • 核心源码:starburst_trino_usage.py
  • 集成测试:test_starburst_trino_usage.py
  • 连接器能力注册表:datahub.json
  • 相关连接器:Trino 元数据摄取可参考 trino_pre.md 与 trino_recipe.yml,两者配合可同时获得 Trino 的元数据、血缘与使用统计。

【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

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

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

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

立即咨询