ARTICLE DETAIL

资讯详情

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

DataHub DataFlow 与 DataJob 实体管理实战:用 Python SDK 构建数据处理管线元数据

DataHub DataFlow 与 DataJob 实体管理实战:用 Python SDK 构建数据处理管线元数据 DataHub DataFlow 与 DataJob 实体管理实战用 Python SDK 构建数据处理管线元数据【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub本篇技术指南围绕 DataHub 中的 DataFlow数据处理管线与 DataJob管线内的作业两大实体讲解如何通过官方 Python SDKdatahub.sdk创建、读取、更新与查询这两类实体并结合仓库源码剖析其 URN 结构、Aspect 存储原理与描述字段优先级等底层机制。读完本文你将能够用可复现的代码为 Airflow DAG、Spark 任务等数据处理流程建立完整的元数据模型并为其附加血缘lineage、标签、归属与描述信息。为什么需要 DataFlow 与 DataJob在一个数据生态系统中数据从采集ingestion到转换transformation再到存储storage往往要经过多个处理阶段。DataFlow 与 DataJob 正是用于描述这一过程的实体DataFlow代表一条完整的数据处理管线pipeline例如一个 Airflow DAG、一个 Spark 批处理应用DataJob代表管线内部的一个可执行作业单元例如 Airflow 中的一个 task、Spark 作业中的某个 stage。通过将二者建模为独立实体用户可以统一定义、管理和监控数据在各个环节中的流转把谁在什么时候、用什么脚本、把哪些数据变成了哪些数据完整记录到 DataHub 的元数据图Metadata Graph中从而支撑血缘分析、影响分析和数据治理。官方对二者的精确定义可见于 SDK 源码注释dataflow.py 将 DataFlow 描述为dataflow 表示一个数据集合datajob.py 则将其描述为数据管线的一个可执行单元例如 Airflow task 或 Spark job。本指南目标创建一个 DataFlow创建一个与 DataFlow 关联的 DataJob。前置条件本教程要求你已经部署 DataHub Quickstart 并导入示例数据详细步骤请参阅 DataHub Quickstart Guide。此外教程代码基于 Python SDK需要确保环境中已安装acryl-datahubmetadata-ingestion 包并已通过datahub init命令配置好 GMS 服务地址与访问令牌。SDK 客户端通过环境变量DATAHUB_GMS_URL、DATAHUB_GMS_TOKEN或~/.datahubenv文件读取连接信息其实现见 main_client.py 中DataHubClient.from_env()的注释说明你也可以直接使用DataHubClient(server..., token...)在代码中指定服务端与令牌。创建 DataFlow以下代码创建了一个名为example_dataflow的 Airflow 数据管线并为其附加了两个标签tag。完整脚本见 dataflow_create.pyfrom datahub.metadata.urns import TagUrn from datahub.sdk import DataFlow, DataHubClient client DataHubClient.from_env() dataflow DataFlow( nameexample_dataflow, platformairflow, descriptionairflow pipeline for production, tags[TagUrn(nameproduction), TagUrn(namedata_engineering)], ) client.entities.upsert(dataflow)DataFlow 构造参数详解对照 SDK 构造函数见 dataflow.pyDataFlow的完整参数如下参数类型必填说明namestr是管线标识同时构成 URN 中的flow_id字段platformstr是所属平台如airflow、spark构成 URN 的orchestrator字段platform_instancestr否平台实例标识如PROD、PROD-US-EAST用于区分多环境/多集群envstr否运行环境默认DEFAULT_ENVPROD会归一化为FabricType枚举并驱动环境Environment搜索过滤display_namestr否展示名未设置时回退为namedescriptionstr否管线描述external_urlstr否外部文档/源码链接如 Airflow DAG 的页面地址custom_propertiesDict[str, str]否自定义键值属性如调度表达式、SLA 等created/last_modifieddatetime否创建/修改时间戳subtypestr否子类型如ETLowners/links/tags/terms/domain各种输入类型否归属人、链接、标签、术语表术语、域parent_container容器 URN否父级容器structured_properties字典否结构化属性structured properties从源码可见构造函数内部会调用DataFlowUrn.create_from_ids(orchestratorplatform, flow_idname, envenv, platform_instanceplatform_instance)生成 URN并默认写入DataFlowInfoClass这一 Aspectenv参数会经过_normalize_env()处理——如果传入值不是合法的FabricType枚举Aspect 中的env字段将被置为None对应的环境搜索过滤将失效见 dataflow.py。最后调用client.entities.upsert(dataflow)执行写入底层会将实体转换为 MCPMetadata Change Proposal并发送给 GMS见 entity_client.py。创建 DataJobDataJob 必须与一个 DataFlow 关联。创建 DataJob 有两种方式直接传入一个DataFlow对象DataJob 将继承该 Flow 的 platform 与 platform instance传入DataFlowUrn及其平台实例。两种方式在源码中是对应DataJob.__init__的两个分支若未提供flow对象则要求必须提供flow_urn否则抛出ValueError见 datajob.py。方式一通过 DataFlow 对象创建完整脚本见 datajob_create_full.pyfrom datahub.metadata.urns import DatasetUrn, TagUrn from datahub.sdk import DataFlow, DataHubClient, DataJob client DataHubClient.from_env() # datajob will inherit the platform and platform instance from the flow dataflow DataFlow( platformairflow, nameexample_dag, platform_instancePROD, descriptionexample dataflow, tags[TagUrn(nametag1), TagUrn(nametag2)], ) datajob DataJob( nameexample_datajob, flowdataflow, inlets[ DatasetUrn(platformhdfs, namedataset1, envPROD), ], outlets[ DatasetUrn(platformhdfs, namedataset2, envPROD), ], ) client.entities.upsert(dataflow) client.entities.upsert(datajob)这里inlets输入数据集与outlets输出数据集共同构成了 DataJob 的数据集级血缘dataset1经过该作业转换为dataset2。从源码看inlets/outlets会被写入DataJobInputOutputClassAspect 的inputDatasets/outputDatasets字段见 datajob.py并最终渲染为 DataHub 血缘图中的上下游关系。方式二通过 DataFlow URN 创建当你只想关联一条已存在的管线例如它已由其他流程创建时可以只提供 URN。完整脚本见 datajob_create_with_flow_urn.pyfrom datahub.metadata.urns import DataFlowUrn, DatasetUrn from datahub.sdk import DataHubClient, DataJob client DataHubClient.from_env() # datajob will inherit the platform and platform instance from the flow datajob DataJob( nameexample_datajob, flow_urnDataFlowUrn( orchestratorairflow, flow_idexample_dag, clusterPROD, ), platform_instancePROD, inlets[ DatasetUrn(platformhdfs, namedataset1, envPROD), ], outlets[ DatasetUrn(platformhdfs, namedataset2, envPROD), ], ) client.entities.upsert(datajob)注意当通过 URN 创建时如果flow_urn.flow_id已经带有platform_instance.前缀例如PROD.example_dagSDK 会自动剥离该前缀只保留真实的 flow 名称见 datajob.py。DataJob 构造参数补充除name、flow/flow_urn、inlets/outlets外DataJob还支持display_name、description、external_url、custom_properties、owners、tags、terms、domain、subtype、fine_grained_lineages等参数见 datajob.py。其中fine_grained_lineages用于声明字段级列级血缘。另一个值得注意的细节DataJob 的DataJobInfoClassAspect 默认携带typemodels.AzkabanJobTypeClass.COMMAND作为默认作业类型见 datajob.py。DataFlow 与 DataJob 的 URN 结构理解 URN 是正确读写实体的前提。DataFlow 与 DataJob 的 URN 均为嵌套结构DataFlow URNurn:li:dataFlow:(orchestrator,flow_id,cluster)orchestrator编排平台如airflow、sparkflow_id管线 ID若配置了platform_instance则会带上PROD.前缀如PROD.example_dagcluster集群/环境标识如PROD、prod、dev。DataJob URNurn:li:dataJob:(dataFlowUrn,job_id)DataJob URN 内嵌完整的 DataFlow URN再拼接作业 ID。URN 的解析与序列化均有单元测试覆盖test_data_flow_urn.py 验证了urn:li:dataFlow:(airflow,def,prod)的解析test_data_job_urn.py 验证了urn:li:dataJob:(urn:li:dataFlow:(airflow,flow_id,prod),job_id)的解析以及get_orchestrator_name()、get_flow_id()、get_env()、get_data_flow_urn()、get_job_id()等访问器。读取 DataFlow以下代码根据 URN 读取 DataFlow 实体并打印关键属性。完整脚本见 dataflow_read.pyfrom datahub.sdk import DataFlowUrn, DataHubClient client DataHubClient.from_env() # Or get this from the UI (share - copy urn) and use DataFlowUrn.from_string(...) dataflow_urn DataFlowUrn( orchestratorairflow, flow_idexample_dataflow, clusterPROD ) dataflow_entity client.entities.get(dataflow_urn) print(DataFlow name:, dataflow_entity.name) print(DataFlow platform:, dataflow_entity.platform) print(DataFlow description:, dataflow_entity.description)示例输出 DataFlow name: example_dataflow DataFlow platform: urn:li:dataPlatform:airflow DataFlow description: airflow pipeline for production注意platform打印出来的是一段完整的urn:li:dataPlatform:airflowURN。DataFlowUrn既支持用关键字构造也支持通过 UI 的 Share → Copy URN 复制后调用DataFlowUrn.from_string(...)解析。描述字段的读取优先级从 dataflow.py 的实现可以看到description属性并非直接读取某个单一字段而是按如下优先级合并EditableDataFlowPropertiesClass可编辑属性 Aspect中的descriptionDataFlowInfoClass信息 Aspect中的description。前者是 UI 与 SDK 编辑入口写入的编辑覆盖层后者是 ingestion 流程写入的原始信息。同理set_description()也会根据当前是否为 ingestion 上下文选择写入哪个 Aspect若处于 ingestion 归因is_ingestion_attribution()中会写入DataFlowInfoClass并在覆盖非 ingestion 描述时发出IngestionAttributionWarning告警见 dataflow.py。DataJob 的描述读取与写入逻辑完全对称见 datajob.py。读取 DataJob以下代码读取 DataJob 并打印其名称、所属 Flow 的 URN 与描述。完整脚本见 datajob_read.pyfrom datahub.sdk import DataFlowUrn, DataHubClient, DataJobUrn client DataHubClient.from_env() # Or get this from the UI (share - copy urn) and use DataJobUrn.from_string(...) # The flow_id carries the flows platform instance as a prefix (PROD.), which is # how DataFlow(platform_instance...) builds its urn. datajob_urn DataJobUrn( flowDataFlowUrn( orchestratorairflow, flow_idPROD.example_dag, clusterPROD ), job_idexample_datajob, ) datajob_entity client.entities.get(datajob_urn) print(DataJob name:, datajob_entity.name) print(DataJob Flow URN:, datajob_entity.flow_urn) print(DataJob description:, datajob_entity.description)示例输出 DataJob name: example_datajob DataJob Flow URN: urn:li:dataFlow:(airflow,PROD.example_dag,PROD) DataJob description: example datajob代码注释中的细节值得注意flow_id携带了 Flow 的平台实例前缀PROD.这正是DataFlow(platform_instance...)构建 URN 的方式。DataJob实体的flow_urn属性由self.urn.get_data_flow_urn()从嵌套 URN 中反解得到见 datajob.py。进阶实战血缘、标签、归属与生命周期管理构建多 Job 管线inlets/outlets 串联血缘当多个 DataJob 属于同一个 Flow且前一个 Job 的输出是后一个 Job 的输入时通过inlets/outlets即可自动串联出完整的管线级血缘。参考 dataflow_with_datajobs.py# metadata-ingestion/examples/library/dataflow_with_datajobs.py from datahub.metadata.urns import DatasetUrn from datahub.sdk import DataFlow, DataHubClient, DataJob client DataHubClient.from_env() # Create the parent DataFlow dataflow DataFlow( platformairflow, namecustomer_360_pipeline, descriptionEnd-to-end pipeline for building customer 360 view, envPROD, ) # Create DataJobs that belong to this flow extract_job DataJob( nameextract_customer_data, flowdataflow, descriptionExtracts customer data from operational databases, outlets[ DatasetUrn(platformsnowflake, namestaging.customers_raw, envPROD), ], ) transform_job DataJob( nametransform_customer_data, flowdataflow, descriptionTransforms and enriches customer data, inlets[ DatasetUrn(platformsnowflake, namestaging.customers_raw, envPROD), ], outlets[ DatasetUrn( platformsnowflake, nameanalytics.customers_enriched, envPROD ), ], ) load_job DataJob( nameload_customer_360, flowdataflow, descriptionLoads final customer 360 view, inlets[ DatasetUrn( platformsnowflake, nameanalytics.customers_enriched, envPROD ), ], outlets[ DatasetUrn(platformsnowflake, nameprod.customer_360, envPROD), ], ) # Upsert all entities client.entities.upsert(dataflow) client.entities.upsert(extract_job) client.entities.upsert(transform_job) client.entities.upsert(load_job) print(fCreated DataFlow: {dataflow.urn}) print(f - Job 1: {extract_job.urn}) print(f - Job 2: {transform_job.urn}) print(f - Job 3: {load_job.urn})这是一个非常贴合真实生产场景的客户 360 视图管线extract → transform → load 三个作业通过数据集首尾相连DataHub 会自动基于DataJobInputOutputClass建立 Job 与 Dataset 之间的血缘边。为 DataFlow 添加完整元数据参考 dataflow_comprehensive.py可以为 Flow 一次性附加调度信息、SLA、归属人、标签、术语表术语与域# metadata-ingestion/examples/library/dataflow_comprehensive.py from datetime import datetime, timezone from datahub.metadata.urns import CorpGroupUrn, CorpUserUrn, GlossaryTermUrn, TagUrn from datahub.sdk import DataFlow, DataHubClient client DataHubClient.from_env() dataflow DataFlow( platformairflow, namedaily_sales_aggregation, display_nameDaily Sales Aggregation Pipeline, platform_instancePROD-US-EAST, envPROD, descriptionAggregates daily sales data from multiple sources and updates reporting tables, external_urlhttps://airflow.company.com/dags/daily_sales_aggregation, custom_properties{ team: analytics, schedule: 0 2 * * *, sla_hours: 4, priority: high, }, createddatetime(2024, 1, 15, tzinfotimezone.utc), last_modifieddatetime.now(timezone.utc), subtypeETL, owners[ CorpUserUrn(jdoe), CorpGroupUrn(data-engineering), ], tags[ TagUrn(nameproduction), TagUrn(namesales), TagUrn(namecritical), ], terms[ GlossaryTermUrn(Classification.Confidential), ], domainurn:li:domain:sales, ) client.entities.upsert(dataflow) print(fCreated DataFlow: {dataflow.urn}) print(fDisplay Name: {dataflow.display_name}) print(fDescription: {dataflow.description}) print(fExternal URL: {dataflow.external_url}) print(fCustom Properties: {dataflow.custom_properties})更新描述与添加归属人DataFlow/DataJob 实体对象支持读-改-写模式先用client.entities.get()取出实体再调用 setter 方法修改最后upsert保存。更新描述参考 dataflow_update_description.py核心是dataflow.set_description(...)后调用client.entities.upsert(dataflow)添加归属人参考 dataflow_add_ownership.py使用dataflow.add_owner((CorpUserUrn(alice), DATAOWNER))、dataflow.add_owner((CorpGroupUrn(analytics-team), DATAOWNER))追加个人与群组归属。通过 Patch 更新 DataJob 血缘对于增量更新SDK 还提供了DataJobPatchBuilder支持以 Patch 方式追加输入/输出数据集、字段级血缘与 Job 间上下游且set_fine_grained_lineages()会整体替换全部字段级血缘。参考 datajob_add_lineage_patch.pyfrom datahub.emitter.mce_builder import ( make_data_job_urn, make_dataset_urn, make_schema_field_urn, ) from datahub.ingestion.graph.client import DataHubGraph, DataHubGraphConfig from datahub.metadata.schema_classes import ( FineGrainedLineageClass as FineGrainedLineage, FineGrainedLineageDownstreamTypeClass as FineGrainedLineageDownstreamType, FineGrainedLineageUpstreamTypeClass as FineGrainedLineageUpstreamType, ) from datahub.specific.datajob import DataJobPatchBuilder datahub_client DataHubGraph(DataHubGraphConfig(serverhttp://localhost:8080)) datajob_urn make_data_job_urn( orchestratorairflow, flow_iddag_abc, job_idtask_456 ) patch_builder DataJobPatchBuilder(datajob_urn) patch_builder.add_input_dataset( make_dataset_urn(platformkafka, nameSampleKafkaDataset, envPROD) ) patch_builder.add_output_dataset( make_dataset_urn(platformhive, nameSampleHiveDataset, envPROD) ) patch_builder.add_input_datajob( make_data_job_urn(orchestratorairflow, flow_iddag_abc, job_idtask_123) ) patch_builder.add_input_dataset_field( make_schema_field_urn( parent_urnmake_dataset_urn( platformhive, namefct_users_deleted, envPROD ), field_pathuser_id, ) ) patch_builder.add_output_dataset_field( make_schema_field_urn( parent_urnmake_dataset_urn( platformhive, namefct_users_created, envPROD ), field_pathuser_id, ) ) lineage1 FineGrainedLineage( upstreamTypeFineGrainedLineageUpstreamType.FIELD_SET, upstreams[ make_schema_field_urn(make_dataset_urn(postgres, raw_data.users), user_id) ], downstreamTypeFineGrainedLineageDownstreamType.FIELD, downstreams[ make_schema_field_urn( make_dataset_urn(postgres, analytics.user_metrics), user_id, ) ], transformOperationIDENTITY, confidenceScore1.0, ) patch_builder.add_fine_grained_lineage(lineage1) patch_builder.set_fine_grained_lineages([lineage1]) patch_mcps patch_builder.build() for patch_mcp in patch_mcps: datahub_client.emit(patch_mcp)删除 DataFlow参考 dataflow_delete.pySDK 支持软删除与硬删除两种模式from datahub.metadata.urns import DataFlowUrn from datahub.sdk import DataHubClient client DataHubClient.from_env() dataflow_urn DataFlowUrn(airflow, old_pipeline, dev) # Soft delete the DataFlow (marks as removed but retains metadata) client.entities.delete(dataflow_urn, hardFalse) print(fSoft deleted DataFlow: {dataflow_urn}) # To hard delete (permanently remove): # client.entities.delete(dataflow_urn, hardTrue)搜索 DataFlow还可以通过搜索接口按名称或平台过滤 Flow。参考 dataflow_search.pyfrom datahub.ingestion.graph.client import DatahubClientConfig, DataHubGraph from datahub.sdk import DataHubClient from datahub.sdk.search_filters import FilterDsl as F graph DataHubGraph(DatahubClientConfig(serverhttp://localhost:8080)) # Search for DataFlows by name using get_urns_by_filter results list( graph.get_urns_by_filter( entity_types[dataFlow], querysales, ) ) print(fFound {len(results)} DataFlows matching sales:\n) for urn in results[:10]: # Show first 10 print(fURN: {urn}) # Search with filters using FilterDsl client DataHubClient.from_env() filtered_results list( client.search.get_urns( filterF.and_( F.entity_type(dataFlow), F.platform(airflow), ) ) ) print(f\nFound {len(filtered_results)} Airflow DataFlows)底层原理小结SDK 如何读写 DataFlow / DataJob综合源码可以总结出 SDK 读写这两类实体的一致模式URN 驱动身份每个实体以 URN 为唯一标识DataFlowUrn三元组orchestrator, flow_id, cluster与嵌套的DataJobUrn是全部读写操作的地基Aspect 承载属性核心属性分散在若干 Aspect 中——DataFlowInfoClass/EditableDataFlowPropertiesClass承载 Flow 信息DataJobInfoClass/EditableDataJobPropertiesClass/DataJobInputOutputClass承载 Job 信息、输入输出数据集与字段级血缘描述字段优先级description属性优先返回可编辑 AspectUI/SDK 写入层其次返回信息 Aspectingestion 写入层upsert 幂等写入client.entities.upsert()将实体序列化为 MCP 提交到 GMS实体已存在时会发出覆盖告警见 entity_client.pybrowse path 自动继承创建 DataJob 时SDK 会从所属 Flow 继承并扩展BrowsePathsV2Class使 Job 在浏览路径上挂在 Flow 之下见 datajob.py。如需进一步探索可继续阅读仓库中的以下资源教程文档dataflow-datajob.mdSDK 实体实现dataflow.py、datajob.py客户端实现main_client.py、entity_client.pyURN 单元测试test_data_flow_urn.py、test_data_job_urn.py全部示例脚本metadata-ingestion/examples/library/目录下的dataflow_*.py与datajob_*.py【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表