DataHub SQL Queries 采集源实战:从 JSONL 查询日志解析血缘与操作元数据
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
sql-queries是 DataHub 元数据采集框架(metadata-ingestion)中用于生产环境血缘摄取的核心模块:它读取一个换行分隔的 JSON(JSONL)文件中记录的 SQL 查询,通过 SQL 解析引擎自动生成表级与字段级血缘(Lineage),并可进一步产出查询实体(Query)、使用统计与操作事件(Operation)。本文以仓库中的 sql-queries_pre.md 文档为主线,结合 sql_queries.py 源码与 集成测试用例 的输入样例,完整讲解该采集源的输入格式、配置项、运行方式与底层实现原理。读完本文,你将能够把任意平台导出的 SQL 查询日志(本地文件或 S3 对象)转化为 DataHub 中可检索、可追溯的血缘与操作元数据。
模块定位与核心能力
根据 sql-queries_pre.md 的说明,sql-queries模块的职责是从 SQL 查询中摄取元数据到 DataHub,它面向生产级摄取工作流(production ingestion workflows),模块特有的能力在文档中单独说明。
从源码实现看,该模块在 DataHub 的 source 注册表中以sql-queries为平台 id 注册,并声明了三项能力(见 sql_queries.py 的装饰器声明):
| 能力 | 支持状态 | 说明 |
|---|---|---|
| LINEAGE_COARSE(粗粒度血缘) | 支持 | 从 SQL 解析出的表级上下游关系 |
| LINEAGE_FINE(细粒度血缘) | 支持 | 从 SQL 解析出的字段级血缘 |
| OPERATION_CAPTURE(操作捕获) | 支持 | 从非 SELECT 查询(INSERT/UPDATE/CREATE 等)解析出操作事件 |
该 source 在框架中的支持状态为GA(General Availability,正式可用)。它的工作方式非常直接:读取一个包含 SQL 查询的换行分隔 JSON 文件,解析这些查询以生成血缘。与 Snowflake、BigQuery 等直接连接数据源采集 usage 的 source 不同,sql-queries不直接连接任何数据库,而是消费已经导出的查询日志,因此它尤其适合离线血缘回填、无法直连数据库的环境以及多平台查询日志统一入湖等场景。
前置条件
官方文档在 Prerequisites 一节 明确了运行前的三项前置要求:
- 网络连通性:确保采集环境能访问查询日志所在位置(本地文件系统或 S3)以及 DataHub GMS;
- 有效的认证凭据:配置可用的 DataHub 访问凭据(
datahub_api下的 token 或用户名密码); - 元数据 API 的读权限:该模块需要读取 DataHub 中已存在的 SchemaMetadata(表结构)等元数据来辅助 SQL 解析。
源码进一步印证了这一点:SqlQueriesSource.__init__在构造时会强制校验ctx.graph(即 DataHub API 客户端)非空,否则直接抛出ValueError("SqlQueriesSource needs a datahub_api from which to pull schema metadata")(见 sql_queries.py#L183-L187)。也就是说,该采集源必须在配置了datahub_api的前提下运行,它依赖 DataHub 作为 schema 解析的知识来源。
输入格式:JSONL 查询文件详解
sql-queries的输入是一个换行分隔的 JSON 文件(NDJSON/JSONL),每一行是一个 JSON 对象,对应一条查询记录。每行的字段由源码中的QueryEntry模型定义(见 sql_queries.py#L471-L477):
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
query | str | 是 | SQL 查询文本,血缘与操作解析的核心输入 |
timestamp | datetime(可空) | 否 | 查询执行时间;支持 Unix 秒级时间戳或任意可被parse_user_datetime解析的日期格式 |
user | CorpUserUrn(可空) | 否 | 执行查询的用户;可为 DataHub corpuser URN 字符串,空字符串会被视为"无执行者"而不会导致该行失败 |
downstream_tables | List[DatasetUrn] | 否 | 显式血缘中的下游表列表 |
upstream_tables | List[DatasetUrn] | 否 | 显式血缘中的上游表列表 |
session_id | str(可空) | 否 | 会话标识,用于跨查询维护临时表映射 |
仓库中的 basic.jsonl 给出了最典型的三字段输入样例:
{"query": "SELECT * FROM snowflake.db.users", "timestamp": 1609459200, "user": "john.doe"} {"query": "INSERT INTO snowflake.db.orders SELECT user_id, product_id, order_date FROM snowflake.db.temp_orders", "timestamp": 1609459260, "user": "jane.smith"} {"query": "CREATE VIEW snowflake.db.user_summary AS SELECT u.id, u.name, COUNT(o.id) as order_count FROM snowflake.db.users u LEFT JOIN snowflake.db.orders o ON u.id = o.user_id GROUP BY u.id, u.name", "timestamp": 1609459320, "user": "admin"} {"query": "UPDATE snowflake.db.users SET last_login = CURRENT_TIMESTAMP WHERE id IN (SELECT DISTINCT user_id FROM snowflake.db.sessions WHERE session_date >= '2021-01-01')", "timestamp": 1609459380, "user": "system"}从上面的样例可以看到,sql-queries能处理SELECT、INSERT ... SELECT、CREATE VIEW、UPDATE、CREATE TABLE AS等多种语句形态,它们分别对应血缘、查询实体、操作事件等不同产出的解析路径。
逐行容错解析
源码中_parse_lines的实现(见 sql_queries.py#L365-L403)体现了面向生产设计的容错策略:
- 空行自动跳过;
- 每一行独立使用
json.loads(line, strict=False)解析,单行格式错误不会中断整个摄取,而是计入num_entries_failed并记录 warning 后继续; - 每处理 1000 行输出一次进度日志,方便观察长文件的处理节奏;
- 若所有行都解析失败,
_report_run_health会将其上报为failure(而非常规 warning),使流水线以非零退出码结束,避免"看似成功实则空跑"的假绿(见 sql_queries.py#L269-L299)。
显式血缘(Explicit Lineage)输入
除了让解析器从 SQL 文本中自动推断血缘,输入行还支持通过upstream_tables/downstream_tables直接声明血缘关系。仓库中的 explicit-lineage.jsonl 展示了这种用法:
{"query": "INSERT INTO snowflake.db.orders SELECT user_id, product_id, order_date FROM snowflake.db.temp_orders", "timestamp": 1609459260, "user": "jane.smith", "upstream_tables": ["snowflake.db.users"], "downstream_tables": ["snowflake.db.audit_log"]}源码中的处理逻辑(见 sql_queries.py#L418-L453)如下:
- 若某行同时提供
upstream_tables和downstream_tables,则走KnownQueryLineageInfo路径,完全信任文件中的显式血缘,不再对该行做 SQL 解析; - 若只提供了其中一侧(如只有上游没有下游),则记录日志说明"部分血缘缺失,回退到 SQL 解析",并按普通查询交给解析聚合器处理;
- 表名字符串会被转换为
make_dataset_urn_with_platform_instance生成的 Dataset URN,因此表名会遵循配置中的platform、platform_instance与env; upstream_tables/downstream_tables必须是列表,传入裸字符串会直接报错(源码特意注释了这一点,防止字符串被逐字符拆开伪造出血缘);列表中的非法条目会被忽略并计入num_invalid_table_entries,同时给出 warning。
配置项全解
SqlQueriesSourceConfig(见 sql_queries.py#L75-L148)继承自PlatformInstanceConfigMixin、EnvConfigMixin和IncrementalLineageConfigMixin,即它同时支持 DataHub 通用的platform_instance(平台实例)、env(环境,默认 PROD)以及增量血缘相关配置。模块自身的配置项如下:
| 配置项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
query_file | str | (必填) | 待摄取的查询文件路径,支持本地路径与s3://URI |
platform | str | (必填) | 生成元数据时使用的平台标识,例如snowflake、bigquery,该值决定 Dataset URN 中的 platform 部分 |
usage | BaseUsageConfig | 默认实例 | 生成使用统计时的配置(如start_time、end_time、bucket_duration) |
use_schema_resolver | bool | true | 是否从 DataHub 读取 SchemaMetadata 辅助 SQL 解析;仅测试时可关闭 |
default_db | str | None | 未限定库名的表默认归属的数据库 |
default_schema | str | None | 未限定 schema 名的表默认归属的 schema |
override_dialect | str | None | 强制指定 SQL 方言,覆盖自动方言检测 |
temp_table_patterns | List[str] | [] | 临时表正则模式列表,用于在血缘摄取中过滤临时表 |
aws_config | AwsConnectionConfig | None | S3 访问配置,当query_file为s3://URI 时必填 |
关键配置项说明
query_file与 S3 支持。query_file既可以是本地文件,也可以是 S3 对象 URI。源码通过is_s3_uri判断,并在模型校验器validate_s3_config中强制要求:当query_file以s3://开头时,必须同时提供aws_config,否则配置阶段直接报错(见 sql_queries.py#L136-L142)。读取 S3 文件时使用smart_open配合aws_config.get_s3_client()建立流式读取(见 sql_queries.py#L301-L322),因此即便查询文件很大也能以流式逐行处理,内存占用可控。
temp_table_patterns(临时表过滤)。这是血缘质量的关键配置。某些平台(如 Athena)没有原生临时表概念,但会用命名约定来"模拟"临时表(例如temp_前缀、_temp后缀)。若不加以过滤,这些中间表会大量污染血缘图。源码规定:
- 模式使用起始锚定匹配(
re.match,与AllowDenyPattern一致),例如temp_能匹配任何以temp_开头的表; - 编译时统一忽略大小写(
re.IGNORECASE); - 模式本身会在配置校验阶段被
re.compile验证,非法正则直接报错(见 sql_queries.py#L124-L134); - 命中临时表模式时,
is_temp_table回调会返回 True 并计入num_temp_table_matches,解析聚合器据此在血缘图中滤除这些表(见 sql_queries.py#L455-L468)。
仓库测试 temp-table-patterns.yml 中的用法示例:
temp_table_patterns: ["^temp_.*", "^tmp_.*", ".*_temp$"]use_schema_resolver(Schema 解析器)。开启后(默认开启),source 会构造SchemaResolver,从 DataHub 拉取已注册表的 SchemaMetadata 作为 SQL 解析的上下文。由于真实世界的 SQL 往往使用未限定的表名、别名、大小写变体,拥有列级 schema 信息能显著提升解析准确率。该配置在源码中被标记为HiddenFromDocs,属于面向测试的开关,生产环境建议保持默认开启。
default_db/default_schema。当查询日志中的表名未带库名或 schema 限定时,这两个配置提供兜底的解析上下文。例如查询只写了SELECT * FROM orders,配合default_db: snowflake、default_schema: db可解析为snowflake.db.orders。
override_dialect。sql-queries底层使用 SQLGlot 解析引擎,通常会自动检测方言;但当自动检测在混合方言日志中表现不稳定时,可通过该配置显式指定方言(如snowflake、bigquery、spark)。
usage(使用统计配置)。通过嵌套的BaseUsageConfig控制使用统计的生成窗口与聚合粒度。仓库测试 basic.yml 中的典型配置:
usage: start_time: 2021-01-01T00:00:00Z end_time: 2021-01-02T00:00:00Z bucket_duration: DAY已移除的配置项
源码通过pydantic_removed_field标记了enable_lazy_schema_loading这一历史配置,它已在2026 年 8 月被正式移除(见 sql_queries.py#L120-L122)。如果你在旧版 recipe 中看到该配置,升级后应直接删除,否则会收到配置校验告警。
完整的 Recipe 示例
以下 recipe 来自仓库集成测试 basic.yml,展示了一个可运行的完整配置骨架:
source: type: sql-queries config: query_file: ./input/basic.jsonl platform: snowflake use_schema_resolver: false usage: start_time: 2021-01-01T00:00:00Z end_time: 2021-01-02T00:00:00Z bucket_duration: DAY sink: type: file config: filename: ./output.json datahub_api: server: http://localhost:8080实际生产环境中,一般建议:
- 将
use_schema_resolver保持为true(默认),并正确配置datahub_api.server指向 DataHub GMS; - 若查询日志表名未完全限定,补充
default_db与default_schema; - 若平台存在临时表命名约定,配置
temp_table_patterns过滤中间表; - 查询文件存放在 S3 时,补充
aws_config(region、凭据等)并将query_file写为s3://bucket/path/queries.jsonl。
运行命令为标准的 DataHub ingestion CLI:
datahub ingest -c recipe.yml底层实现原理:解析聚合器与 Schema 解析
sql-queries的核心设计是"解析与产出分离":摄取阶段把每条查询喂给SqlParsingAggregator,产出阶段再由聚合器统一生成各类元数据工作单元(workunit)。
两阶段摄取流程
get_workunits_internal(见 sql_queries.py#L242-L267)将一次运行划分为两个报告阶段:
- QUERIES_EXTRACTION(查询抽取):逐行读取查询文件,把每条查询加入聚合器。单条查询加入失败(如解析异常)不会中止运行,而是计入
num_queries_aggregator_failures并记录 warning,但系统性错误(内存溢出、进程中断、网络认证类异常)会被重新抛出以终止任务; - LINEAGE_EXTRACTION(血缘抽取):调用
aggregator.gen_metadata()生成所有血缘、查询、使用统计与操作元数据,并通过auto_workunit包装为统一的工作单元流输出。
SqlParsingAggregator 的初始化参数
聚合器的构造(见 sql_queries.py#L207-L225)揭示了模块的能力开关,这些参数当前在源码中标记为 TODO 待配置化,但已经全部启用:
| 参数 | 值 | 作用 |
|---|---|---|
generate_lineage | true | 生成血缘 |
generate_queries | true | 生成 Query 实体 |
generate_query_subject_fields | true | 生成查询主题字段 |
generate_query_usage_statistics | true | 发布 SELECT 查询实体(仅当开启时才会为 SELECT 发布 Query 实体,否则只发布写操作类查询) |
generate_usage_statistics | true | 生成使用统计 |
generate_operations | true | 生成操作事件(从非 SELECT 查询解析) |
eager_graph_load | false | 不从 DataHub 预加载全量 schema,按需惰性加载 |
is_temp_table | 配置了temp_table_patterns时为回调 | 临时表判定回调 |
format_queries | false | 是否格式化查询文本 |
会话与临时表追踪
源码实现注释明确指出,模块通过session_id跨查询维护临时表映射:同一会话中,CREATE TEMP TABLE之后的查询引用该临时表时,解析器能正确追踪其真实上游。这一能力由集成测试 session-temp-tables.jsonl 对应的session-temp-tables用例验证,其 golden 文件位于 session-temp-tables.json。
增量血缘与补丁格式
由于配置类继承了IncrementalLineageConfigMixin,该 source 还支持增量血缘:每次运行可以只处理新增的查询窗口,避免全量重算。同时测试目录中的patch-lineage用例表明,血缘可以采用patch 格式(MCP Patch)输出,而非全量覆盖式的 upstreamLineage,这为频繁增量摄取提供了更低开销的更新方式。
测试验证与产出样例
仓库在 tests/integration/sql_queries 目录下提供了完整的集成测试套件,覆盖了该模块的九类典型场景,可直接作为理解模块行为的"活文档":
| 测试用例 | 验证点 |
|---|---|
basic | 基础血缘解析(SELECT/INSERT/CREATE VIEW/UPDATE/CTAS 混合场景) |
basic-with-schema-resolver | 开启 Schema Resolver 后的解析差异 |
session-temp-tables | 同一 session 内临时表血缘的正确追踪 |
query-deduplication | 重复查询的去重 |
explicit-lineage | 文件内显式声明血缘的信任路径 |
hex-origin | 十六进制 origin 相关场景 |
patch-lineage | 血缘以 patch 格式输出 |
lazy-schema-resolver | 惰性 schema 加载 |
temp-table-patterns | 临时表正则过滤对血缘图的净化效果 |
测试运行方式(见 test_sql_queries.py)值得注意:它通过 docker-compose 启动一个MockServer 模拟 DataHub,将datahub_api.server指向临时端口(通过环境变量SQL_QUERIES_MOCK_PORT注入),随后用Pipeline.create(recipe)真实执行摄取,最后将输出与 golden 文件比对。这说明该模块的端到端链路(读文件 → 解析 → 聚合 → 产出 → 上报)完全可以在无真实 DataHub 实例的情况下被验证。
各用例的输入文件(.jsonl+.ymlrecipe)与期望输出(golden/*.json)一一对应,例如 basic.jsonl 的期望血缘输出见 basic.json,适合读者对照学习每类 SQL 语句最终会生成哪些 MCP 元数据。
运行状态与排障
SqlQueriesSourceReport(见 sql_queries.py#L151-L161)提供了完善的运行观测指标,运行结束后可通过datahub ingest的输出或报告 API 查看:
| 指标 | 含义 |
|---|---|
num_entries_processed | 成功解析的输入行数 |
num_entries_failed | 解析失败被跳过的行数(单行坏数据不影响整体) |
num_queries_processed | 成功加入解析聚合器的查询数 |
num_queries_aggregator_failures | 加入聚合器失败的查询数 |
num_invalid_table_entries | 显式血缘中被忽略的非法表引用数 |
num_temp_table_matches | 命中临时表模式被过滤的表数 |
sql_aggregator/schema_resolver_report | 聚合器与 schema 解析器的子报告 |
排障时可关注以下已知边界行为:
- 所有行解析失败:报告会以 failure 级别提示"check the file format (expected newline-delimited JSON)",通常意味着文件并非 JSONL 格式或编码异常;
- 所有查询聚合失败:报告会提示可能是认证、连通性、配置等系统性问题的信号;
- 空输入:文件存在但没有查询条目时,会给出 warning 而非 failure;
- 空字符串 user:很多平台导出的查询日志会把系统/后台查询的 user 记为
"",源码将其视同"无执行者"处理,不会丢弃整行,也不会因构造 URN 失败而中断。
总结
sql-queries是 DataHub 血缘体系中对"查询日志型"元数据源的标准答案:它用一个 JSONL 文件加一个 recipe 即可完成从 SQL 日志到血缘、查询实体、使用统计与操作事件的完整摄取,天然适合离线回填与多云日志统一治理。其源码(sql_queries.py)与集成测试(test_sql_queries.py)为生产落地提供了完整的配置参考、容错策略与验证样例——在接入自己的查询日志前,建议先用仓库中的样例 JSONL 文件跑通端到端链路,再逐步替换为真实数据源并调优temp_table_patterns、default_db/default_schema等血缘质量相关配置。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考