拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Microsoft Fabric Data Factory 元数据接入 DataHub:概念映射、权限配置与血缘解析实战

Microsoft Fabric Data Factory 元数据接入 DataHub概念映射、权限配置与血缘解析实战【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub导读fabric-data-factory是 DataHub 官方元数据接入源ingestion source用于将 Microsoft Fabric Data Factory 中的工作区Workspace、数据管道Data Pipeline与活动Activity等编排实体同步到 DataHub并基于 Copy 与 InvokePipeline 活动解析数据集级血缘、抽取管道与活动的执行历史。读完本文你将掌握该连接器的概念映射模型、四种 Azure 认证方式的配置、权限与 Fabric 管理员设置、血缘解析与跨 Recipe 对齐platform_instance_map的方法以及执行历史、多租户隔离、状态化删除检测等高级用法。一、连接器概览与概念映射Microsoft Fabric Data Factory 是 Fabric 平台内的云数据集成服务DataHub 通过fabric-data-factory连接器将其中的编排实体建模为 DataHub 实体。官方集成目录将其归类为「ETL/ELT / Pull 类」连接器见 integrations_catalog.json。连接器覆盖的实体范围包括工作区、数据管道、活动Copy、Lookup、Spark 等、管道运行与活动运行同时捕获表级血缘与状态化删除检测。核心概念映射关系如下Fabric Data Factory 概念DataHub 实体说明Workspace工作区Container子类型Fabric Workspace顶层组织单元Data Pipeline数据管道DataFlow包含活动的编排管道Activity活动DataJob管道内单个任务Copy、Lookup、Spark 等Pipeline Run管道运行DataProcessInstance管道运行的执行记录Activity Run活动运行DataProcessInstance管道内单个活动的执行记录Connection连接解析为外部 Dataset用于解析到外部平台数据集的血缘层级结构连接器产出的实体层级如下Platform (fabric-data-factory) └── Workspace (Container) └── Data Pipeline (DataFlow) └── Activity (DataJob) ├── Pipeline Run (DataProcessInstance) └── Activity Run (DataProcessInstance)从源码看该层级关系由 source.py 中的_emit_pipelines实现工作区通过build_workspace_container生成Container每个管道以DataFlow实体发布携带pipeline_id、workspace_id自定义属性与 Fabric Web UI 外部链接管道内的活动则作为子DataJob发布。核心能力工作区建模为容器数据管道建模为 DataFlowDataHub 实体类型活动建模为 DataJob从 Copy 与 InvokePipeline 活动提取数据集级血缘管道与活动执行历史建模为 DataProcessInstance通过platform_instance_map实现跨 Recipe 血缘与外部导入的数据集正确相连支持对工作区与管道按名称模式pattern过滤支持状态化摄取stateful ingestion以移除陈旧实体支持四种认证方式服务主体Service Principal、托管身份Managed Identity、Azure CLI、DefaultAzureCredential。二、快速开始配置认证—— 按「前置条件」小节配置 Azure 凭据启用 API 访问—— 若使用服务主体或托管身份需由 Fabric 管理员启用服务主体 API 访问授予权限—— 将你的身份添加为工作区Contributor读取管道定义与血缘所需配置 Recipe—— 以 fabric-data-factory_recipe.yml 为模板运行摄取—— 执行datahub ingest -c fabric-data-factory_recipe.yml。三、前置条件权限、Fabric 管理员设置与认证3.1 所需权限连接器要求在每个工作区上具备Contributor角色。没有 Contributor 则无法获取管道定义仅有 Reader 角色时连接器可以列出工作区与管道但不会提取管道活动、活动运行详情与血缘。委派认证以用户身份若使用委派认证如 Azure CLI登录用户已有的 Fabric 权限直接生效。连接器需要以下委派范围scopeWorkspace.Read.All或Workspace.ReadWrite.All—— 用于列出工作区与条目Item.ReadWrite.All或DataPipeline.ReadWrite.All—— 用于获取条目定义Get Item Definition、列出条目连接List Item Connections与查询活动运行Query Activity RunsItem.Read.All不足够获取定义与连接Item.Read.All或DataPipeline.Read.All—— 足够用于列出条目作业实例List Item Job Instances即执行历史。Azure CLI 令牌默认已包含必要的 Fabric API scope。服务主体与托管身份认证服务主体SP与托管身份MI默认不继承任何权限需要启用 API 访问Fabric 管理员需启用服务主体租户设置见下文「Fabric 管理员设置」授予工作区访问权将 SP 或 MI 添加为每个待摄取工作区的Contributor。3.2 Fabric 管理员设置⚠️ 对于服务主体与托管身份认证Fabric 管理员必须在 Fabric 管理门户中为服务主体启用 API 访问。否则即使工作区权限分配正确API 调用也会以 401 失败。截至 2025 年年中Microsoft 将原来的单一租户设置拆分为两个独立设置应按如下方式配置进入 Fabric 管理门户 租户设置Tenant settings在开发者设置Developer settings下启用适用的设置项服务主体可调用 Fabric 公共 APIService principals can call Fabric public APIs—— 控制受 Fabric 权限模型保护的 CRUD API 访问如读取工作区与条目。新租户自 2025 年 8 月起默认启用。服务主体可创建工作区、连接与部署管道Service principals can create workspaces, connections, and deployment pipelines—— 控制不受 Fabric 权限保护的全局 API 访问默认禁用仅在需要时启用。将访问权限限制到专用安全组仅包含需要 API 访问的服务主体。这是推荐做法。 若你的租户较旧仍能看到旧版单一设置服务主体可使用 Fabric APIService principals can use Fabric APIs请直接启用它系统会自动迁移为上述两个新设置。 租户设置变更传播最长可能需要15 分钟。启用后立即收到 401 错误时请稍等重试。3.3 认证方式连接器通过共享的credential配置块支持四种认证方式全部基于 Azure 的TokenCredential接口。统一认证实现在 azure_auth.py 的AzureCredentialConfig中并校验不同认证方式所需的必填字段。服务主体生产环境推荐在 Microsoft Entra ID 中注册应用记下client_id、client_secret、tenant_id然后确保 Fabric 管理员已启用服务主体 API 访问见上在 Entra ID 中创建安全组并将服务主体加入其中将安全组添加为每个目标工作区的ContributorContributor 角色可访问管道定义与条目连接用于血缘解析。credential: authentication_method: service_principal client_id: ${AZURE_CLIENT_ID} client_secret: ${AZURE_CLIENT_SECRET} tenant_id: ${AZURE_TENANT_ID}使用该方式时三个字段均必填。源码中 validate_credentials 会校验缺失项并抛出明确错误。托管身份适用于 Azure 托管的部署在 Azure VM、AKS、App Service 或其他支持托管身份的 Azure 计算资源上运行时使用。托管身份必须添加为 Fabric 工作区 Contributor。同时仍需 Fabric 管理员启用上文「Fabric 管理员设置」中的租户设置——该设置虽以「服务主体」命名但同样管辖托管身份的 API 访问。# 系统分配托管身份无需额外配置 credential: authentication_method: managed_identity对于用户分配托管身份需提供 client IDcredential: authentication_method: managed_identity managed_identity_client_id: your-managed-identity-client-idAzure CLI本地开发与测试使用本地az login会话的凭据。登录用户已有的 Fabric 权限直接生效除工作区访问权限外无需额外设置。credential: authentication_method: cli开始摄取前先运行az login。无浏览器的远程服务器可使用az login设备码流程。DefaultAzureCredential灵活的自动检测使用 Azure 的 DefaultAzureCredential 链按顺序尝试多种凭据来源环境变量、工作负载身份workload identity、托管身份、共享令牌缓存、Azure CLI、Azure PowerShell、Azure Developer CLI 等。credential: authentication_method: default可以从链中排除特定凭据来源以加快检测速度或在混合环境中避免意外认证credential: authentication_method: default exclude_cli_credential: true # 跳过 Azure CLI生产环境推荐 exclude_environment_credential: false exclude_managed_identity_credential: false对应地azure_auth.py 在default模式下将这些排除项直接透传给DefaultAzureCredential。3.4 设置步骤汇总从上文选择一种认证方式并配置credential块若使用服务主体或托管身份确保 Fabric 管理员启用了相应开发者设置创建安全组、加入你的身份并在目标工作区授予Contributor若使用 Azure CLI运行az login远程服务器可加--use-device-code配置摄取 Recipe 及可选的工作区、管道过滤。四、完整 Recipe 配置详解以下为连接器的完整示例 Recipe即 fabric-data-factory_recipe.yml# Example recipe for Fabric Data Factory source # See README.md for full configuration options source: type: fabric-data-factory config: # Authentication (using service principal) credential: authentication_method: service_principal client_id: ${AZURE_CLIENT_ID} client_secret: ${AZURE_CLIENT_SECRET} tenant_id: ${AZURE_TENANT_ID} # Optional: Filter workspaces by name pattern workspace_pattern: allow: - .* # Allow all workspaces by default deny: [] # Optional: Filter pipelines by name pattern pipeline_pattern: allow: - .* # Allow all pipelines by default deny: [] # Feature flags extract_pipelines: true include_lineage: true include_execution_history: true execution_history_days: 7 # 1-90 days # Optional: Map Fabric connection names to platform instances for accurate lineage # platform_instance_map: # my-snowflake-connection: prod_snowflake # my-bigquery-connection: analytics_project # Optional: Platform instance for this Fabric Data Factory connector # platform_instance: my-fabric-tenant # Environment env: PROD # Optional: Stateful ingestion for stale entity removal # stateful_ingestion: # enabled: true sink: type: datahub-rest config: server: http://localhost:8080 token: ${DATAHUB_GMS_TOKEN}各配置项在 config.py 的FabricDataFactorySourceConfig中有明确定义配置项默认值说明credentialAzureCredentialConfig()Azure 认证配置支持服务主体/托管身份/Azure CLI/自动检测workspace_pattern全部允许按名称正则过滤工作区如allow[prod-.*]、deny[.*-test]pipeline_pattern全部允许按名称正则过滤数据管道作用于所有匹配workspace_pattern的工作区extract_pipelinestrue是否提取数据管道及其活动include_lineagetrue从活动输入/输出提取血缘将 Fabric 连接按类型映射到 DataHub 数据集include_execution_historytrue将管道与活动执行历史提取为 DataProcessInstance含运行状态、时长与参数execution_history_days7提取执行历史的天数1–90仅在include_execution_history为 true 时生效值越大摄取时间越长Fabric API 每个管道最多返回 100 条最近完成运行platform_instance_map{}连接名到 DataHub platform instance 的映射如{my-snowflake-connection: prod_snowflake}用于将血缘解析到已有数据集api_timeout30REST API 调用超时秒范围 1–300stateful_ingestion无状态化摄取配置启用后跟踪已摄取实体并移除 Fabric 中已不存在的实体其中execution_history_days通过 pydantic 的ge1、le90约束取值区间api_timeout约束在 1–300 秒。连接器能力声明从 source.py 的装饰器可以看到连接器的能力声明容器默认启用、平台实例默认启用、粗粒度血缘默认通过 Copy 与 InvokePipeline 活动启用、删除检测通过 stateful_ingestion 可选启用支持状态标记为BETA。五、血缘提取深度解析5.1 哪些活动产生血缘连接器从以下 Fabric 活动类型提取数据集级血缘活动类型血缘行为Copy从输入数据集到输出数据集建立血缘InvokePipeline建立指向子管道的管道到管道血缘血缘默认启用include_lineage: true。5.2 血缘解析如何工作要让血缘正确连接到其他源如 Snowflake、BigQuery摄取的数据集连接器需要将 Fabric 连接解析为 DataHub 平台。第一步自动连接映射连接器自动将 Fabric 连接类型映射到 DataHub 平台例如Snowflake连接映射到snowflake平台。完整映射表见 constants.py 中的FABRIC_CONNECTION_PLATFORM_MAP覆盖 OneLakeLakehouse/Warehouse→fabric-onelake、SQL Server→mssql、MySQL、Oracle、PostgreSQL、Azure Blob/ADLS→abs、BigQuery、Redshift、Snowflake、Hive、Spark、Databricks、Dremio、Salesforce、Vertica、Kusto、Cassandra、Kafka、GCS、HDFS、S3、MongoDB、Presto、dbt、Elasticsearch、Looker、Delta Sharing 等。未映射的连接类型会回退为使用连接类型字符串作为平台名。从 lineage.py 的_resolve_platform实现可见解析顺序先查FABRIC_CONNECTION_PLATFORM_MAP再回退到 ADF 的ADF_LINKED_SERVICE_PLATFORM_MAP例如AzureBlobStorage→abs最后才以连接类型字符串本身兜底并记录report_unmapped_connection_type报告。对应单元测试见 test_lineage.py。第二步平台实例映射跨 Recipe 血缘如果你用其他 DataHub 连接器如 Snowflake、BigQuery摄取同一数据源必须确保platform_instance值一致。用platform_instance_map将 Fabric 连接名映射到其他 Recipe 中使用的平台实例# Fabric Data Factory Recipe source: type: fabric-data-factory config: credential: authentication_method: service_principal client_id: ${AZURE_CLIENT_ID} client_secret: ${AZURE_CLIENT_SECRET} tenant_id: ${AZURE_TENANT_ID} platform_instance_map: # Key: Your Fabric connection name (exact match required) # Value: The platform_instance from your other source recipe snowflake-prod-connection: prod_warehouse bigquery-analytics: analytics_project# Corresponding Snowflake Recipe (platform_instance must match) source: type: snowflake config: platform_instance: prod_warehouse # Must match the value in platform_instance_map # ... other configplatform_instance_map的解析优先级在 lineage.py 中实现优先按连接名查映射表未命中则回退到全局platform_instance再未配置则为空。单元测试 test_lineage.py 覆盖了「映射命中 / 未命中回退全局 / 无映射无全局」三种分支。⚠️ 若platform_instance值不匹配血缘将创建独立的数据集实体而不会连接到已摄取的数据集。5.3 Copy 活动血缘的源码级解析Copy 活动的血缘解析由CopyActivityLineageExtractor完成lineage.py。它从活动的typeProperties.source与typeProperties.sink/destination中读取datasetSettings再按以下顺序解析连接_resolve_connection_and_type见 lineage.pyexternalReferences.connection→ 连接缓存 → 取(display_name, connection_type) 2a.connectionSettings.properties.externalReferences.connection→ 连接缓存 2b.connectionSettings内联的(name, properties.type)linkedService的(name, properties.type)。对于 Fabric 原生类型Lakehouse、Warehouse解析为与 OneLake 连接器一致的 URN平台fabric-onelake格式{workspace}.{item}.{schema}.{table}无 schema 时默认dbo见 urn_generator.py对于外部平台则基于schema.table或文件路径container/fileSystem/bucketName folderPath fileName构造标准DatasetUrn。主流程 source.py 会对「源端缺失」「汇端缺失」「两端都缺失」分别给出警告或失败报告并在异常时降级为「活动照常发布但不带血缘」保证单个活动的失败不影响整个管道。5.4 InvokePipeline 跨管道血缘InvokePipelineLineageExtractorlineage.py只支持InvokeFabricPipeline这一种操作类型。解析时读取活动的pipelineId子管道与workspaceId缺省用父工作区从活动缓存中查找子管道的根活动无上游依赖、depends_on为空的活动多个根时取第一个在父活动的 DataJob 上挂calls_pipeline_id、calls_workspace_id、child_pipeline_urn、operation_type、child_root_activity等自定义属性并向子活动 DataJob 追加父级上游边_cross_pipeline_edges与_merge_cross_pipeline_info见 source.py。值得注意的是摄取采用两趟two-pass设计第一趟遍历所有工作区与管道将活动缓存到_pipeline_activities_cache并抓取连接缓存_connections_cache第二趟才解析边并发布实体。这保证了跨工作区、跨管道的 InvokePipeline 边解析不依赖处理顺序见 source.py。六、执行历史Execution History管道与活动运行默认作为DataProcessInstance实体提取source: type: fabric-data-factory config: include_execution_history: true # default execution_history_days: 7 # 1-90 days提供的信息包括运行状态、时长、时间戳、调用类型invoke type以及活动级细节错误消息、重试次数。从源码看管道运行与活动运行分别走_emit_pipeline_run与_emit_activity_runsource.py管道运行 DPI 携带run_id、workspace_id、pipeline_id可选invoke_type、failure_reason类型为BATCH_SCHEDULED模板 URN 指向所属管道 DataFlow活动运行 DPI 携带activity_run_id、pipeline_run_id、activity_type可选duration_in_ms、error_message、retry_attempt并通过parent_instance关联到所属管道运行 DPI状态映射管道运行侧COMPLETED→SUCCESS、FAILED→FAILURE、CANCELLED/DEDUPED→SKIPPEDsource.py活动运行侧Succeeded/Failed/Cancelled三态映射source.py每次运行发布 generate/start/end 三类 MCP 事件_emit_dpi_workunits见 source.py活动运行查询使用queryActivityRunsAPIPOST/v1/workspaces/{wsId}/datapipelines/pipelineruns/{jobId}/queryactivityruns见 client.py起始时间下限设为 2015-01-01ACTIVITY_RUNS_MIN_START确保不遗漏任意历史的运行。 Fabric API 每个管道最多返回100 条最近完成的运行。如需更深的运行历史请提高摄取频率。七、多租户与高级配置7.1 何时使用platform_instance从多个环境摄取时用连接器的platform_instance配置来区分独立的 Fabric 租户场景风险解决方案单租户无不需要多租户高—— 名称冲突风险必须使用# Multi-tenant example source: type: fabric-data-factory config: platform_instance: contoso-tenant # Prevents URN collisions⚠️ 不同 Fabric 租户可能存在同名的工作区与管道。使用platform_instance防止实体被覆盖。7.2 URN 格式管道 URN 格式urn:li:dataFlow:(fabric-data-factory,{workspace_id}.{pipeline_id},{env})带platform_instance时urn:li:dataFlow:(fabric-data-factory,{platform_instance}.{workspace_id}.{pipeline_id},{env})URN 生成集中在 urn_generator.py管道 DataFlow 的 flow id 为{workspaceGUID}.{pipelineGUID}活动 DataJob 的 job id 直接使用活动名Fabric 保证同一管道内活动名唯一活动运行 DPI id 为{pipelineRunGUID}.{activityRunGUID}从而携带对所属管道运行的直接引用。7.3 状态化摄取与陈旧实体清理启用stateful_ingestion后连接器会记录已摄取实体清单并在后续运行中移除 Fabric 中已不存在的实体删除检测。配置如下source: type: fabric-data-factory config: stateful_ingestion: enabled: true其配置类继承自StatefulIngestionConfigBase与StatefulStaleMetadataRemovalConfig见 config.py底层由StatefulIngestionSourceBase与StatefulStaleMetadataRemovalHandler提供状态存取与陈旧元数据移除机制。八、限制与已知边界运行历史上限Fabric API 每个管道最多返回 100 条最近完成运行。若execution_history_days覆盖的运行超过该上限只返回最近 100 条。应提高摄取频率以捕获更深历史。不支持 Dataflow Gen2工作区级、带转换逻辑的独立 Dataflow Gen2 条目不会被提取。不支持 CopyJob工作区级独立 CopyJob 条目不会被提取只有管道内嵌的 Copy 活动才产生血缘。无触发器/调度元数据管道触发器与调度信息不会被提取。不支持 ExecutePipelineExecutePipeline活动类型在 Fabric 中标记为旧版legacy不支持跨管道血缘。血缘范围仅 Copy 与 InvokePipeline 活动产生数据集或管道血缘其他活动类型Lookup、Wait、ForEach、Script 等作为 DataJob 摄取但不带数据集级血缘。InvokePipeline 活动操作类型仅支持InvokeFabricPipeline操作类型做跨管道血缘InvokeAdfPipeline、InvokeExternalPipeline不会解析将被跳过对应 lineage.py 的告警路径。基于查询的 Copy 源当 Copy 活动使用sqlReaderQuery或sqlReaderStoredProcedureName而非直接表引用时不提取血缘。无列级血缘连接器仅提取数据集级血缘不提取 Copy 活动 translator 配置中的列到列映射。无 Notebook/SparkJobDefinition 血缘Notebook 与 SparkJobDefinition 活动会作为 DataJob 摄取但血缘不解析。连接解析兜底未映射的连接类型会回退为使用连接类型字符串作为平台名可能与 DataHub 中已有的平台名不一致。请用platform_instance_map显式映射连接名。九、故障排查401/403 错误确认服务主体具备正确的 Fabric API 权限并已添加为工作区成员。结果为空检查workspace_pattern与pipeline_pattern是否把所有条目都过滤掉了。血缘缺失确认include_lineage: true已设置且管道的 Fabric 连接配置正确同时对照「限制」小节中不支持的活动类型与场景进行排查。陈旧实体启用stateful_ingestion自动移除 Fabric 中已不存在的实体。十、源码导读若希望深入理解该连接器的实现可按以下路径阅读当前仓库连接器入口与主流程source.py两趟摄取、实体发布、运行处理配置模型与校验config.py血缘解析器lineage.pyFabric REST 客户端client.py连接类型到平台映射constants.pyURN 生成规范urn_generator.py统一 Azure 认证azure_auth.py单元测试test_lineage.py 与 test_urn_generator.py官方集成目录条目integrations_catalog.json【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门