DataHub Starburst Trino Usage 连接器:基于 Event Logger 摄取 Trino 使用统计的完整指南 📅 发布时间:2026/9/19 22:03:34 👁 浏览次数: 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生产可用能力CapabilityUSAGE_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.pyemail_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与envTop 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返回预置的访问事件驱动完整Pipelinesource 为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 名用户显示为xxxtest.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),仅供参考