
使用 DataHub Python SDK 构建与查询数据血缘Lineage SDK 完整实战指南【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahubDataHub 的 Python SDKacryl-datahub为骨架系统讲解如何通过LineageClient添加表级/列级血缘、从 SQL 自动推断血缘、按多跳hop与结构化过滤条件读取上游/下游血缘并给出基于仓库源码metadata-ingestion/src/datahub/sdk/lineage_client.py与可运行示例metadata-ingestion/examples/library的实现原理与完整可复制代码。读完本文你将能够在 DataHub 实例上以编程方式完成血缘的写入、推断、查询与过滤并理解其底层 GraphQL 调用与 SQL 解析机制。环境准备与连接 DataHub使用 DataHub SDK 前需要先安装acryl-datahub包并完成与 DataHub 实例的连接配置。安装方式与 CLI 安装一致可参考 CLI 安装指南 中关于acryl-datahub的说明。连接 DataHub 实例from datahub.sdk import DataHubClient client DataHubClient(serveryour_server, tokenyour_token)serverDataHub GMS 服务的 URL。本地部署http://localhost:8080云端托管https://your_datahub_url/gmstoken需要从你的 DataHub 实例中生成个人访问令牌Personal Access Token。此外所有示例代码普遍使用DataHubClient.from_env()方式初始化客户端它从环境变量如DATAHUB_GMS_URL、DATAHUB_GMS_TOKEN等读取连接信息更适合在脚本、调度任务或 CI/CD 中使用见示例 lineage_get_basic.py。添加血缘add_lineage()方法用于在两个实体之间定义血缘关系。它在 lineage_client.py 中以重载形式支持多组实体组合内部会根据「上游实体类型 → 下游实体类型」的组合分发到不同的处理器若传入了不支持的组合或对非 Dataset→Dataset 血缘误传column_lineage/transformation_text会抛出SdkUsageError。添加实体级血缘upstream与downstream参数传入需要建立关联的实体 URN支持 Dataset、DataJob、Dashboard、Chart 等实体类型。数据集Dataset之间建立血缘from datahub.metadata.urns import DatasetUrn from datahub.sdk.main_client import DataHubClient client DataHubClient.from_env() upstream_urn DatasetUrn(platformsnowflake, namesales_raw) downstream_urn DatasetUrn(platformsnowflake, namesales_cleaned) client.lineage.add_lineage(upstreamupstream_urn, downstreamdownstream_urn)完整示例见 lineage_dataset_add.py。这里通过DatasetUrn构造器便捷生成标准 URN完整 URN 形如urn:li:dataset:(urn:li:dataPlatform:snowflake,sales_raw,PROD)无需手工拼接字符串。数据任务DataJob之间建立血缘from datahub.metadata.urns import DataFlowUrn, DataJobUrn from datahub.sdk import DataHubClient client DataHubClient.from_env() dataflow_urn DataFlowUrn( orchestratorairflow, flow_iddata_pipeline, clusterPROD ) client.lineage.add_lineage( upstreamDataJobUrn(flowdataflow_urn, job_iddata_job_1), downstreamDataJobUrn(flowdataflow_urn, job_iddata_job_2), )完整示例见 lineage_datajob_to_datajob.py。DataFlowUrn与DataJobUrn分别对应数据流水线如 Airflow DAG与流水线内的单个任务便于在血缘图中表达作业编排关系。[!NOTE] 血缘组合支持范围 支持的实体组合见下文「支持的血缘组合」一节并非所有组合都允许添加列级血缘。添加列级血缘列级血缘Column-Level LineageCLL通过add_lineage()的column_lineage参数启用仅在 Dataset→Dataset 组合下生效。模糊匹配Fuzzy Matchingfrom datahub.metadata.urns import DatasetUrn from datahub.sdk import DataHubClient client DataHubClient.from_env() client.lineage.add_lineage( upstreamDatasetUrn(platformsnowflake, namesales_raw), downstreamDatasetUrn(platformsnowflake, namesales_cleaned), column_lineageTrue, # same as auto_fuzzy, which maps columns based on name similarity )完整示例见 lineage_dataset_column.py。当column_lineageTrue时DataHub 会基于列名相似度自动映射列适用于上下游列名相近但不完全一致的情况例如上游customer_id对应下游CustomerId。其底层实现位于 lineage_client.py 的_get_fuzzy_column_lineage先对列名做归一化处理再通过difflib等相似度算法完成匹配。严格匹配Strict Matchingfrom datahub.metadata.urns import DatasetUrn from datahub.sdk import DataHubClient client DataHubClient.from_env() client.lineage.add_lineage( upstreamDatasetUrn(platformsnowflake, namesales_raw), downstreamDatasetUrn(platformsnowflake, namesales_cleaned), column_lineageauto_strict, )完整示例见 lineage_dataset_column_auto_strict.py。严格模式下上游与下游的列名必须完全一致才能建立映射。从源码看_get_strict_column_lineagelineage_client.py实际采用大小写不敏感的匹配先构建小写化映射再按下游列名逐一匹配因此customer_id与Customer_Id这类仅大小写不同的列名仍可匹配而拼写不同的列名不会建立血缘。自定义映射Custom Mapping当自动匹配无法满足复杂映射关系时可传入字典自定义映射其中键为下游列名值为上游列名列表从而表达一对多等复杂关系from datahub.metadata.urns import DatasetUrn from datahub.sdk import DataHubClient client DataHubClient.from_env() client.lineage.add_lineage( upstreamDatasetUrn(platformsnowflake, namesales_raw), downstreamDatasetUrn(platformsnowflake, namesales_cleaned), # { downstream_column - [upstream_columns] } column_lineage{ id: [id], region: [region, region_id], total_revenue: [revenue], }, )完整示例见 lineage_dataset_column_custom_mapping.py。上例中下游region列同时依赖上游region与region_id两列下游total_revenue由上游revenue计算而来可精确刻画字段级加工逻辑。从 SQL 自动推断血缘infer_lineage_from_sql()可以解析一段 SQL 查询自动确定上游与下游数据集并添加血缘——包括尽可能生成列级血缘同时创建一个携带 SQL 变换逻辑的查询节点Query Nodefrom datahub.sdk.main_client import DataHubClient client DataHubClient.from_env() sql_query CREATE TABLE sales_summary AS SELECT p.product_name, c.customer_segment, SUM(s.quantity) as total_quantity, SUM(s.amount) as total_sales FROM sales s JOIN products p ON s.product_id p.id JOIN customers c ON s.customer_id c.id GROUP BY p.product_name, c.customer_segment # sales_summary will be assumed to be in the default db/schema # e.g. prod_db.public.sales_summary client.lineage.infer_lineage_from_sql( query_textsql_query, platformsnowflake, default_dbprod_db, default_schemapublic, )完整示例见 lineage_dataset_from_sql.py。方法签名lineage_client.py还支持platform_instance、env默认PROD与override_dialect参数参数默认值说明query_text必填待解析的 SQL 查询文本platform必填数据平台标识如snowflake、bigqueryplatform_instanceNone平台实例名用于多实例部署场景envPROD环境标识default_db/default_schemaNone未显式指定库/模式时用于补全表名如示例中sales_summary会被解析为prod_db.public.sales_summaryoverride_dialectNone覆盖方言推断用于方言检测不准的场景从源码看该方法调用create_lineage_sql_parsed_result完成解析内部基于 DataHub SQL Parser其核心依赖 sqlglot解析失败如无输出表会抛出SdkUsageError随后以第一个输出表作为下游遍历所有输入表建立血缘并跳过自引用self-lineage边列级映射会按「输出列 → 对应上游列列表」组织后一并写入。解析器对SELECT、CREATE TABLE AS SELECT、INSERT、UPDATE等语句支持表级与列级血缘支持子查询、CTE、UNION ALL、SELECT *展开等能力但对标量 UDF、UNNEST、Struct 展开、多语句 SQL 脚本等场景支持有限需要在实践中注意。[!NOTE] DataHub SQL Parser 关于 SQL 解析的更多细节请参考 DataHub SQL Parser 文档。携带变换文本添加查询节点当向add_lineage()提供transformation_text时DataHub 会创建一个表示变换逻辑的查询节点Query Node用于追踪数据集之间的加工过程from datahub.metadata.urns import DatasetUrn from datahub.sdk.main_client import DataHubClient client DataHubClient.from_env() upstream_urn DatasetUrn(platformsnowflake, nameupstream_table) downstream_urn DatasetUrn(platformsnowflake, namedownstream_table) transformation_text from pyspark.sql import SparkSession spark SparkSession.builder.appName(HighValueFilter).getOrCreate() df spark.read.table(customers) high_value df.filter(lifetime_value 10000) high_value.write.saveAsTable(high_value_customers) client.lineage.add_lineage( upstreamupstream_urn, downstreamdownstream_urn, transformation_texttransformation_text, column_lineage{id: [id, customer_id]}, ) # by passing the transformation_text, the query node will be created with the table level lineage. # transformation_text can be any transformation logic e.g. a spark job, an airflow DAG, python script, etc. # if you have a SQL query, we recommend using add_dataset_lineage_from_sql instead. # note that transformation_text itself will not create a column level lineage.完整示例见 lineage_dataset_add_with_query_node.py。transformation_text可以是任意变换逻辑Python 脚本、Airflow DAG 代码或任何描述上游数据如何变换为下游数据的代码。[!NOTE] 重要提示 仅提供transformation_text不会自动生成列级血缘必须显式指定column_lineage参数才能启用列级血缘。 如果变换逻辑本身就是 SQL 查询建议改用infer_lineage_from_sql()它会自动解析查询并同时添加列级血缘见 docs/api/tutorials/lineage.md。获取血缘get_lineage()用于检索指定实体的血缘信息。获取实体级血缘获取数据集的上游血缘默认只返回紧邻的1 跳上游实体from datahub.metadata.urns import DatasetUrn from datahub.sdk.main_client import DataHubClient client DataHubClient.from_env() downstream_lineage client.lineage.get_lineage( source_urnDatasetUrn(platformsnowflake, namesales_summary), directiondownstream, ) print(downstream_lineage)完整示例见 lineage_get_basic.py。direction支持upstream上游与downstream下游两个方向。跨多跳获取上下游血缘要获取超过一跳的上下游实体可设置max_hops参数从而在血缘图中按指定跳数遍历from datahub.metadata.urns import DatasetUrn from datahub.sdk.main_client import DataHubClient client DataHubClient.from_env() downstream_lineage client.lineage.get_lineage( source_urnDatasetUrn(platformsnowflake, namesales_summary), directiondownstream, max_hops2, ) print(downstream_lineage)完整示例见 lineage_get_with_hops.py。从源码lineage_client.py可以看到实现细节当max_hops 2时会把degree过滤条件设为[1, 2]等具体数值当max_hops 2时则遍历完整血缘图并将 degree 过滤条件设为[1, 2, 3]结果数量默认上限为 500 条可通过count参数调整同时会打印一条日志提示。[!NOTE] 使用 max_hops 若max_hops大于 2SDK 将尝试遍历完整血缘图并通过count限制返回结果数量。返回类型get_lineage()返回LineageResult对象列表定义见 lineage_client.pyresults [ LineageResult( urnurn:li:dataset:(urn:li:dataPlatform:snowflake,table_2,PROD), typeDATASET, hops1, directiondownstream, platformsnowflake, nametable_2, # name of the entity paths[] # Only populated for column-level lineage ) ]各字段含义urn为实体 URNtype为实体类型hops为距离源实体的跳数direction为查询方向platform为数据平台name为实体名称paths仅在列级血缘查询时填充。获取列级血缘按列获取下游血缘通过source_column参数指定列名即可返回包含该列的血缘路径from datahub.metadata.urns import DatasetUrn from datahub.sdk.main_client import DataHubClient client DataHubClient.from_env() # Get column lineage for the entire flow # you can pass source_urn and source_column to get lineage for a specific column # alternatively, you can pass schemaFieldUrn to source_urn. # e.g. source_urnurn:li:schemaField:(urn:li:dataset:(urn:li:dataPlatform:snowflake,downstream_table),id) downstream_column_lineage client.lineage.get_lineage( source_urnDatasetUrn(platformsnowflake, namesales_summary), source_columnid, directiondownstream, ) print(downstream_column_lineage)完整示例见 lineage_column_get.py。底层实现lineage_client.py会把source_column组装为SchemaFieldUrn并通过searchFlags.groupingSpec以SCHEMA_FIELD为基本实体进行分组聚合搜索。也可以直接把SchemaFieldUrn作为source_urn传入from datahub.sdk.main_client import DataHubClient client DataHubClient.from_env() # Get column lineage for the entire flow results client.lineage.get_lineage( source_urnurn:li:schemaField:(urn:li:dataset:(urn:li:dataPlatform:snowflake,sales_summary,PROD),id), directiondownstream, ) print(list(results))完整示例见 lineage_column_get_from_schemafield.py。返回类型列级血缘的返回类型与实体级一致但paths字段会被填充为列血缘路径LineagePath对象列表定义见 lineage_client.pyresults [ LineageResult( urnurn:li:dataset:(urn:li:dataPlatform:snowflake,table_2,PROD), typeDATASET, hops1, directiondownstream, platformsnowflake, nametable_2, # name of the entity paths[ LineagePath( urnurn:li:schemaField:(urn:li:dataset:(urn:li:dataPlatform:snowflake,table_1,PROD),col1), column_namecol1, # name of the column entity_nametable_1, # name of the entity that contains the column ), LineagePath( urnurn:li:schemaField:(urn:li:dataset:(urn:li:dataPlatform:snowflake,table_2,PROD),col4), column_namecol4, # name of the column entity_nametable_2, # name of the entity that contains the column ) ] # Only populated for column-level lineage ) ]LineagePath中的urn为列SchemaFieldURNcolumn_name为列名entity_name为包含该列的实体名称。结果的解释方法见下文「解读列级血缘结果」。过滤血缘结果get_lineage()支持按平台platform、实体类型type、域domain、环境environment等进行结构化过滤过滤语法与 Search SDK 一致FilterDslfrom datahub.sdk.main_client import DataHubClient from datahub.sdk.search_filters import FilterDsl as F client DataHubClient.from_env() # get upstream snowflake production datasets. results client.lineage.get_lineage( source_urnurn:li:dataset:(platform,sales_agg,PROD), directionupstream, filterF.and_(F.platform(snowflake), F.entity_type(dataset), F.env(PROD)), ) print(results)完整示例见 lineage_get_with_filter.py。F.and_表示各条件取交集F.platform、F.entity_type、F.env分别对应平台、实体类型与环境过滤。SDK 内部会把Filter编译为 GraphQL 的orFilters结构并下发给服务端执行lineage_client.py同时自动叠加degree过滤以匹配max_hops。更多可用过滤器请参考 Search SDK 文档基于过滤器的搜索。Lineage SDK 参考支持的血缘组合血缘 API 支持以下实体组合与 lineage_client.py 中注册的处理器一一对应上游实体下游实体DatasetDatasetDatasetDataJobDataJobDataJobDataJobDatasetDatasetDashboardChartDashboardDashboardDashboardDatasetChartℹ️ 列级血缘与携带变换文本创建查询节点仅支持Dataset → Dataset血缘。若对其他组合传入column_lineage或transformation_textSDK 会抛出SdkUsageError见 lineage_client.py。列级血缘选项对于 Dataset→Dataset 血缘add_lineage()的column_lineage参数支持多种取值取值说明False禁用列级血缘默认True启用列级血缘并自动映射等价于auto_fuzzyauto_fuzzy启用列级血缘采用模糊匹配适合列名相近的场景auto_strict启用列级血缘采用严格匹配要求列名完全一致列映射字典下游列名到上游列名列表的映射[!NOTE]auto_fuzzy与auto_strict的区别auto_fuzzy基于列名相似度自动匹配容忍命名约定上的差异。例如以下列对会被视为匹配user_id→userIdcustomer_id→CustomerIdauto_strict要求上下游列名精确一致。例如上游的customer_id必须与下游的customer_id完全一致才能建立血缘实现上为大小写不敏感匹配。解读列级血缘结果查询列级血缘时结果中的paths展示了列在数据集之间的关联方式每条路径是从源列到目标列的列 URN 序列。例如假设三张表之间存在如下血缘链路table_1.col1 ── table_2.col4 / table_2.col5 ── table_3.col7max_hops1示例 client.lineage.get_lineage( source_urnurn:li:dataset:(urn:li:dataPlatform:snowflake,table_1,PROD), source_columncol1, directiondownstream, max_hops1 )返回[ { urn: ...table_2..., hops: 1, paths: [ [...table_1.col1, ...table_2.col4], [...table_1.col1, ...table_2.col5] ] } ]max_hops2示例 client.lineage.get_lineage( source_urnurn:li:dataset:(urn:li:dataPlatform:snowflake,table_1,PROD), source_columncol1, directiondownstream, max_hops2 )返回[ { urn: ...table_2..., hops: 1, paths: [ [...table_1.col1, ...table_2.col4], [...table_1.col1, ...table_2.col5] ] }, { urn: ...table_3..., hops: 2, paths: [ [...table_1.col1, ...table_2.col4, ...table_3.col7] ] } ]可以看到一跳查询只返回直达列映射两跳查询则把table_2作为中间节点继续向下游传播路径中完整记录了col1 → col4 → col7的传播链便于进行影响分析impact analysis与数据追溯。替代方案Lineage GraphQL API虽然我们通常推荐使用 Python SDK 处理血缘但 DataHub 同样提供 GraphQL API 来添加与检索血缘。用 GraphQL 在数据集之间添加血缘mutation updateLineage { updateLineage( input: { edgesToAdd: [ { downstreamUrn: urn:li:dataset:(urn:li:dataPlatform:hive,logging_events,PROD) upstreamUrn: urn:li:dataset:(urn:li:dataPlatform:hive,fct_users_deleted,PROD) } ] edgesToRemove: [] } ) }edgesToAdd声明新增的血缘边上游→下游edgesToRemove用于移除已有边。用 GraphQL 获取下游血缘query scrollAcrossLineage { scrollAcrossLineage( input: { query: * urn: urn:li:dataset:(urn:li:dataPlatform:hive,logging_events,PROD) count: 10 direction: DOWNSTREAM orFilters: [ { and: [ { condition: EQUAL negated: false field: degree values: [1, 2, 3] } ] } ] } ) { searchResults { degree entity { urn type } } } }scrollAcrossLineage支持分页滚动遍历血缘结果degree过滤字段与 SDK 中的max_hops语义对应1、2、3分别表示一跳、两跳与多跳。用 GraphQL 按时间过滤血缘通过lineageFlags按血缘边的最后更新时间过滤query searchAcrossLineage { searchAcrossLineage( input: { query: * urn: urn:li:dataset:(urn:li:dataPlatform:snowflake,analytics.orders,PROD) count: 10 direction: UPSTREAM orFilters: [{ and: [{ field: degree, values: [1] }] }] lineageFlags: { startTimeMillis: 1625097600000 endTimeMillis: 1627776000000 } } ) { searchResults { entity { urn type } degree } } }该查询只返回在 2021 年 7 月 1 日至 8 月 1 日之间最后更新过的上游血缘边适用于「该时间段内哪些数据依赖关系发生了变化」这类时间敏感的分析。常见问题FAQ能否获取列级血缘可以——对于 Dataset→Dataset 血缘add_lineage()与get_lineage()均支持列级血缘的写入与读取。能否传入 SQL 查询自动获得血缘可以——使用infer_lineage_from_sql()解析查询并提取表级与列级血缘同时生成携带 SQL 变换逻辑的查询节点。检索血缘时能否使用过滤器可以——get_lineage()通过FilterDsl接收结构化过滤器用法与 Search SDK 一致可组合平台、实体类型、环境等多个条件。血缘关系存到哪里、如何生效add_lineage()底层通过 MetadataChangeProposalWrapperMCP以 PATCH 方式如DatasetPatchBuilder、DataJobPatchBuilder等把血缘边写入 DataHub GMSget_lineage()则组装 GraphQL 查询底层对应scrollAcrossLineage等接口从服务端拉取结果因此写入后即可在同一实例内查询验证。延伸阅读Lineage SDK 完整参考源码级 API 文档涵盖全部方法与类型定义DataHub SQL Parser 详解SQL 解析能力、支持/不支持语法与限制Search SDK 文档FilterDsl过滤语法详解更多可运行示例本指南全部示例位于 metadata-ingestion/examples/library文件名前缀为lineage_的脚本【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考