ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

DataHub SQL Queries 采集源实战:从 JSONL 查询日志解析血缘与操作元数据

DataHub SQL Queries 采集源实战:从 JSONL 查询日志解析血缘与操作元数据 DataHub SQL Queries 采集源实战从 JSONL 查询日志解析血缘与操作元数据【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahubsql-queries是 DataHub 元数据采集框架metadata-ingestion中用于生产环境血缘摄取的核心模块它读取一个换行分隔的 JSONJSONL文件中记录的 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 在框架中的支持状态为GAGeneral 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字段类型必填说明querystr是SQL 查询文本血缘与操作解析的核心输入timestampdatetime可空否查询执行时间支持 Unix 秒级时间戳或任意可被parse_user_datetime解析的日期格式userCorpUserUrn可空否执行查询的用户可为 DataHub corpuser URN 字符串空字符串会被视为无执行者而不会导致该行失败downstream_tablesList[DatasetUrn]否显式血缘中的下游表列表upstream_tablesList[DatasetUrn]否显式血缘中的上游表列表session_idstr可空否会话标识用于跨查询维护临时表映射仓库中的 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, strictFalse)解析单行格式错误不会中断整个摄取而是计入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与envupstream_tables/downstream_tables必须是列表传入裸字符串会直接报错源码特意注释了这一点防止字符串被逐字符拆开伪造出血缘列表中的非法条目会被忽略并计入num_invalid_table_entries同时给出 warning。配置项全解SqlQueriesSourceConfig见 sql_queries.py#L75-L148继承自PlatformInstanceConfigMixin、EnvConfigMixin和IncrementalLineageConfigMixin即它同时支持 DataHub 通用的platform_instance平台实例、env环境默认 PROD以及增量血缘相关配置。模块自身的配置项如下配置项类型默认值说明query_filestr必填待摄取的查询文件路径支持本地路径与s3://URIplatformstr必填生成元数据时使用的平台标识例如snowflake、bigquery该值决定 Dataset URN 中的 platform 部分usageBaseUsageConfig默认实例生成使用统计时的配置如start_time、end_time、bucket_durationuse_schema_resolverbooltrue是否从 DataHub 读取 SchemaMetadata 辅助 SQL 解析仅测试时可关闭default_dbstrNone未限定库名的表默认归属的数据库default_schemastrNone未限定 schema 名的表默认归属的 schemaoverride_dialectstrNone强制指定 SQL 方言覆盖自动方言检测temp_table_patternsList[str][]临时表正则模式列表用于在血缘摄取中过滤临时表aws_configAwsConnectionConfigNoneS3 访问配置当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_resolverSchema 解析器。开启后默认开启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_configregion、凭据等并将query_file写为s3://bucket/path/queries.jsonl。运行命令为标准的 DataHub ingestion CLIdatahub 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_lineagetrue生成血缘generate_queriestrue生成 Query 实体generate_query_subject_fieldstrue生成查询主题字段generate_query_usage_statisticstrue发布 SELECT 查询实体仅当开启时才会为 SELECT 发布 Query 实体否则只发布写操作类查询generate_usage_statisticstrue生成使用统计generate_operationstrue生成操作事件从非 SELECT 查询解析eager_graph_loadfalse不从 DataHub 预加载全量 schema按需惰性加载is_temp_table配置了temp_table_patterns时为回调临时表判定回调format_queriesfalse是否格式化查询文本会话与临时表追踪源码实现注释明确指出模块通过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),仅供参考
返回列表