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 |
| 平台 ID | trino |
| 支持状态 | GA(生产可用) |
| 能力(Capability) | USAGE_STATS,默认启用 |
前提条件(Prerequisites)
官方文档(starburst-trino-usage_pre.md)明确列出运行摄取前必须满足的三个前提:
- 网络连通性:摄取节点必须能够访问 Trino 集群(或托管审计库的数据库)的网络端点,即配方中的
host_port可达; - 有效认证凭据:提供具有查询审计表权限的用户名与密码;
- 元数据 API 读权限:账号需要对 Event Logger 暴露的审计元数据(
completed_queries表及相关字段)具备只读权限。
此外,由于实现依赖 Starburst Event Logger,运行环境需要满足:
- Trino 侧已开启并配置好Event Logger,能够持续写入
completed_queries审计表; - 审计表位于某个可被该连接器直接查询的 catalog/schema 下,通过配方中的
audit_catalog与audit_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_port | 是 | Trino 服务地址(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(继承自TrinoConfig与BaseUsageConfig)一一对应(starburst_trino_usage.py):
email_domain:将被追加到用户名的邮箱域名;audit_catalog/audit_schema:审计表所在位置;database:从中获取使用统计的 catalog 名。
继承自基类的进阶参数
除了配方中显式列出的字段,TrinoUsageConfig还继承了TrinoConfig与BaseUsageConfig的大量参数,可在需要时按需追加,包括但不限于:
- 时间窗口:
start_time/end_time,控制查询completed_queries的时间范围(默认值遵循BaseUsageConfig,源码在拼装 SQL 时通过strftime("%Y-%m-%d %H:%M:%S.%f %Z")格式化后注入查询); - 聚合粒度:
bucket_duration,决定按多长时间窗口聚合访问事件(源码通过get_time_bucket对starttime取桶); - 用户过滤:
user_email_pattern,在写入读事件时用于匹配用户邮箱; - 查询采样:
top_n_queries、include_top_n_queries、format_sql_queries、queries_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、权限与元数据暴露情况,结合源码可梳理出如下约束:
- 仅统计成功的 SELECT 查询:查询 SQL 硬编码了
query_type = 'SELECT'且query_state = 'FINISHED'的条件,DDL、失败查询、进行中查询都不会被计入; - 依赖 Starburst Event Logger:必须存在填充了
completed_queries表的审计 catalog/schema,非 Starburst 发行版或未开启 Event Logger 的 Trino 无法使用本模块; - 仅聚合单一 catalog:通过
database配置项过滤,accessed_metadata中 catalog 不等于该值的访问会被直接跳过; - 系统查询被忽略:catalog 以
$system@开头的查询会从统计中排除(源码第 250-257 行); - 用户身份需要可解析为邮箱:源码中对无邮箱的用户会拼接
{user}@{email_domain},无法识别用户名的场景会回退为unknown@{email_domain}; 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 MetadataWorkUnit1. 审计查询:读取 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模型(字段含catalogName、schema、table、columns、connectorInfo)来承载它。_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_instance与env;Top N 查询相关的top_n_queries、format_sql_queries、include_top_n_queries、queries_character_limit等参数在此处生效。
5. 报告与可观测性
TrinoUsageReport在SourceReport基础上增加了num_joined_access_events_skipped计数。摄取出错或部分事件被跳过时,可通过该计数与日志(如Field accessed_metadata is empty. Skipping ....、The username parameter is missing. Skipping ....)快速定位是数据缺失还是解析失败。
测试验证:如何确认模块行为
仓库提供了针对该模块的集成测试 test_starburst_trino_usage.py,包含两类用例:
- 配置解析测试(
test_trino_usage_config):验证TrinoUsageConfig能正确解析host_port、database、username、password、email_domain、audit_catalog、audit_schema、include_views、include_tables等字段,是校验配方字段书写的参考范本; - 摄取端到端测试(
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)是:
- 先验证基础四要素:凭据(Credentials)是否有效、权限(Permissions)是否覆盖审计表、网络连通性(Connectivity)、范围过滤(Scope filters,即
database/ 时间窗口)是否正确; - 再检查摄取日志:针对 source 特有错误信息定位并调整配置。
结合源码,常见的具体问题与对策:
| 现象 | 可能原因 | 处置建议 |
|---|---|---|
日志出现SQL Result is empty,无任何输出 | 时间窗口内无 FINISHED 的 SELECT 查询,或审计表未写入 | 扩大start_time/end_time窗口;确认 Event Logger 已启用 |
| 统计中缺少某 catalog 的数据 | database未设为该 catalog,或访问记录中catalogName与配置不一致 | 核对database与审计记录中的 catalog 名 |
用户显示为xxx@test.com或unknown@... | 用户名非邮箱格式 | 确认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_port、database、audit_catalog、audit_schema、email_domain的配方即可运行;若需更精细的控制(时间窗口、Top N 查询、用户过滤等),可在继承自TrinoConfig与BaseUsageConfig的参数中按需扩展。
延伸阅读
- 模块官方文档: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),仅供参考