ARTICLE DETAIL

资讯详情

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

DataHub 接入 Microsoft Fabric Data Factory:数据血缘提取、执行历史与多租户配置实战指南

DataHub 接入 Microsoft Fabric Data Factory:数据血缘提取、执行历史与多租户配置实战指南 DataHub 接入 Microsoft Fabric Data Factory数据血缘提取、执行历史与多租户配置实战指南【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub导读本文将围绕 DataHub 的fabric-data-factory摄取源Source系统讲解如何把 Microsoft Fabric 数据工厂中的工作区、数据管道Data Pipeline、活动Activity与执行历史同步进 DataHub并重点剖析其数据集级血缘解析机制、跨 Recipe 血缘打通方案、多租户场景下的 URN 隔离策略以及当前版本的能力边界与常见故障排查方法。读完本文你将能独立完成 Fabric Data Factory 的认证配置、Recipe 编写、血缘验证与多租户部署并理解其底层实现原理。一、核心能力总览fabric-data-factory模块从 Microsoft Fabric Data Factory 摄取元数据到 DataHub覆盖数据管道与编排类实体工作区、数据管道、活动并捕获表级血缘与状态化删除检测stale entity removal。模块入口位于 source.py配置模型位于 config.py。概念映射Fabric 中的对象与 DataHub 实体并非一一同名官方概念映射如下Fabric Data Factory 概念DataHub 实体说明Workspace工作区Containersubtype:Fabric Workspace顶层组织单元Data Pipeline数据管道DataFlow包含活动的编排管道Activity活动DataJob管道内的单个任务Copy、Lookup、Spark 等Pipeline RunDataProcessInstance管道运行的执行记录Activity RunDataProcessInstance管道内单个活动的执行记录Connection连接解析为外部平台 Dataset用于将血缘解析到外部平台上的数据集对应层级结构如下Platform (fabric-data-factory) └── Workspace (Container) └── Data Pipeline (DataFlow) └── Activity (DataJob) ├── Pipeline Run (DataProcessInstance) └── Activity Run (DataProcessInstance)能力矩阵速览能力默认值是否需要额外配置容器Workspace → Container开启否平台实例platform_instance关闭多租户时必须配置粗粒度血缘Copy / InvokePipeline开启include_lineage: true跨 Recipe 血缘需platform_instance_map删除检测stale entity removal关闭通过stateful_ingestion开启上述能力声明可以在 source.py 的装饰器capability中看到该连接器当前标注为BETA支持状态SupportStatus.BETA。二、血缘提取哪些活动产生血缘如何解析2.1 产生血缘的活动类型连接器从以下 Fabric 活动类型中提取数据集级血缘活动类型血缘行为Copy从输入数据集到输出数据集创建血缘InvokePipeline创建指向子管道的管道到管道pipeline-to-pipeline血缘血缘默认开启include_lineage: true。除这两种活动外Lookup、Wait、ForEach、Script 等活动类型仅作为DataJob被摄取不产生数据集级血缘。2.2 血缘解析原理两步映射为了让血缘能够正确连接到由其他摄取源如 Snowflake、BigQuery 连接器已经摄入 DataHub 的数据集连接器需要把 Fabric 连接解析到 DataHub 平台。Step 1自动连接映射Automatic Connection Mapping连接器自动将 Fabric 连接类型映射为 DataHub 平台例如Snowflake连接映射为snowflake平台。完整的映射表定义在 constants.py 的FABRIC_CONNECTION_PLATFORM_MAP中覆盖了 50 余种连接类型核心分组如下连接类型Fabric映射平台DataHubLakehouse/Warehouse/FabricSql/DataLake/SqlAnalyticsEndpoint/FabricSqlEndpointMetadatafabric-onelakeCopyJob/FabricDataPipelinesfabric-data-factoryPowerBI/PowerBIDatasets/PowerBIDatamarts/PowerPlatformDataflows等powerbiSQL/AzureSqlMI/Synapse/AzureSynapseWorkspacemssqlMySql/AzureDatabaseForMySQL/MariaDBForPipelinemysqlOracle/AmazonRdsForOracleoraclePostgreSQL/AzurePostgreSQLpostgresAzureBlobs/AzureDataLakeStorageabsGoogleBigQuery/GoogleBigQueryAadbigqueryAmazonRedshiftredshiftSnowflakesnowflakeApacheHive/AzureHivehiveSparksparkDatabricks/DatabricksMultiCloud/AzureDatabricksWorkspacedatabricksSalesforce/SalesforceServiceCloudsalesforceAmazonS3/AmazonS3Compatibles3GoogleCloudStoragegcsAzureDataExplorerkustoConfluentCloudkafkaDeltaSharingdelta-lakeElasticSearch/Looker/Dremio/Vertica/Cassandra/Presto/MongoDB*等对应平台名不支持的连接类型会回退fallback为直接使用连接类型字符串作为平台名例如SapHana→ 平台SapHana这可能导致平台名与 DataHub 中已有平台不匹配详见 lineage.py 的_resolve_platform实现。从源码看映射解析还包含一层 ADF 兼容回退如果连接类型在 Fabric 映射表中找不到会继续尝试ADF_LINKED_SERVICE_PLATFORM_MAPAzure Data Factory 链接服务映射例如AzureBlobStorage、AzureSqlDatabase都找不到时才把连接类型字符串本身当作平台名。Step 2平台实例映射Platform Instance Mapping用于跨 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: Fabric 连接名称需精确匹配 # Value: 其他源 Recipe 中的 platform_instance snowflake-prod-connection: prod_warehouse bigquery-analytics: analytics_project# 对应的 Snowflake Recipeplatform_instance 必须匹配 source: type: snowflake config: platform_instance: prod_warehouse # 必须与 platform_instance_map 中的值一致 # ... 其他配置⚠️关键约束若platform_instance不匹配血缘会创建独立的数据集实体而不会连接到已摄入的既有数据集。从源码实现看lineage.py 的_resolve_platform_instance逻辑是优先用连接名在platform_instance_map中查找映射值找不到则回退到全局的platform_instance配置。2.3 源码级解析链路血缘提取的底层链路见 source.py 与 lineage.pyCopyActivityLineageExtractor读取活动的typeProperties中的source/sink或destination的datasetSettings通过_resolve_connection_and_type解析连接名与连接类型解析顺序为externalReferences.connection→ 查连接缓存connection cacheconnectionSettings.properties.externalReferences.connection→ 查连接缓存connectionSettings.properties.type内联类型linkedService.properties.type链接服务类型对于 Fabric 原生类型Lakehouse / Warehouse生成与 OneLake 连接器一致的 URNfabric-onelake平台格式为{workspaceGUID}.{itemGUID}.{schema}.{table}对于外部平台提取表名支持schema.table或基于location块的文件路径如 container/fileSystem/bucketName folderPath fileName生成标准DatasetUrn。需要说明的是血缘解析依赖在摄取过程中缓存的连接信息连接器会为每个管道调用list_item_connections填充_connections_cache见 source.py。如果连接不在缓存中血缘会解析失败并产生告警。InvokePipeline 血缘由InvokePipelineLineageExtractor处理只有operationType InvokeFabricPipeline时才会解析跨管道血缘其余操作类型InvokeAdfPipeline、InvokeExternalPipeline会被跳过。解析时通过pipelineId/workspaceId找到子管道再定位其根活动depends_on为空的首个活动作为子 DataJob URN详见 lineage.py。三、执行历史DataProcessInstance 摄取管道与活动的运行记录默认以DataProcessInstance实体摄取source: type: fabric-data-factory config: include_execution_history: true # 默认值 execution_history_days: 7 # 取值范围 1-90 天这提供了运行状态run status、时长duration、时间戳、触发方式invoke type以及活动级别的细节包括错误信息error message和重试次数retry attempts。配置参数细节从 config.py 看include_execution_history默认true是否将管道/活动执行历史提取为DataProcessInstance。开启后还能使用真实运行时的参数值从参数化活动中提取血缘execution_history_days默认7ge1, le90提取的执行历史天数仅在include_execution_history为真时生效。取值越大摄取时间越长。Fabric API 的 100 次运行上限:::note Fabric API 每个管道最多返回100 条最近完成的运行记录。如需捕获更久远的历史请提高摄取频率更频繁地运行摄取。 :::也就是说即使execution_history_days覆盖的运行数超过 100也只会返回最近的 100 条。源码级实现执行历史的两级 DPI 生成逻辑位于 source.py管道运行 DPI通过get_pipeline_runs获取回看窗口内的运行跳过NOT_STARTED状态生成DataProcessInstancetype 为BATCH_SCHEDULED附带run_id、workspace_id、pipeline_id、invoke_type、failure_reason等自定义属性并缓存其 URN 供活动运行关联父实例活动运行 DPI通过query_activity_runs查询每个管道运行下的活动运行查询起始时间被设为 2015-01-01 以确保全覆盖生成子级 DPI附带activity_type、duration_in_ms、error_message、retry_attempt等属性并通过parent_instance关联到所属的管道运行 DPI。状态映射同样在源码中可见Fabric 的COMPLETED→SUCCESS、FAILED→FAILURE、CANCELLED/DEDUPED→SKIPPED活动运行的Succeeded/Failed/Cancelled字符串映射到对应的InstanceRunResult。每个 DPI 会按“创建 → 开始事件 → 结束事件”三段发出 MCPMetadata Change Proposal。四、多租户设置使用platform_instance区分租户何时使用platform_instance当从多个环境摄取 Fabric 时使用连接器的platform_instance配置来区分不同的 Fabric 租户场景风险解决方案单租户无不需要多租户高—— 命名冲突风险必须配置# 多租户示例 source: type: fabric-data-factory config: platform_instance: contoso-tenant # 防止 URN 冲突:::warning 不同的 Fabric 租户可能拥有同名的工作区和管道。请使用platform_instance防止实体被互相覆盖entity overwrites。 :::URN 格式管道DataFlowURN 遵循以下格式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其中make_pipeline_flow_urn通过DataFlowUrn.create_from_ids组合 orchestratorfabric-data-factory、flow_id{workspaceGUID}.{pipelineGUID}、env 与 platform_instance。活动 DataJob 的 URN 则基于其所属 DataFlow URN 追加活动名生成Fabric 保证活动名在管道内唯一。五、前提条件与认证配置虽然血缘与执行历史是本文主线但正确的前置配置是这一切可运行的基础。以下内容来自同目录的 fabric-data-factory_pre.md是官方文档的完整前置章节。5.1 快速开始配置认证—— 配置 Azure 凭据见下方认证章节启用 API 访问—— 如使用服务主体SP或托管身份需由 Fabric 管理员启用服务主体 API 访问授权—— 将你的身份添加为工作区Contributor获取管道定义与血缘所必需配置 Recipe—— 以 fabric-data-factory_recipe.yml 为模板运行摄取—— 执行datahub ingest -c fabric-data-factory_recipe.yml。5.2 所需权限连接器要求在每个工作区上拥有Contributor角色否则无法获取管道定义。若只有 Reader 角色连接器能列出工作区与管道但无法提取管道活动、活动运行详情或血缘。委托认证以用户身份所需 delegated 作用域Workspace.Read.All或Workspace.ReadWrite.All—— 列出工作区与项Item.ReadWrite.All或DataPipeline.ReadWrite.All—— 获取项定义、列出项连接、查询活动运行注意Item.Read.All不足以获取定义与连接Item.Read.All或DataPipeline.Read.All—— 足以列出项作业实例执行历史。Azure CLI 令牌默认包含所需的 Fabric API 作用域。服务主体与托管身份认证默认不继承任何权限需要① Fabric 管理员启用服务主体租户设置② 在每个目标工作区将 SP/MI 添加为Contributor。5.3 Fabric 管理员设置:::warning 对于服务主体和托管身份认证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 访问的服务主体的专用安全组中。若为旧租户且仍可见旧式单一设置Service principals can use Fabric APIs启用该设置即可它会被自动迁移为两个新设置。租户设置变更传播最长需要15 分钟若启用后立即遇到 401请稍等重试。5.4 四种认证方式连接器通过共享的credential配置块支持四种认证方式均基于 Azure 的TokenCredential接口。服务主体生产环境推荐credential: authentication_method: service_principal client_id: ${AZURE_CLIENT_ID} client_secret: ${AZURE_CLIENT_SECRET} tenant_id: ${AZURE_TENANT_ID}使用此方式时三个字段均必填。建议在 Entra ID 中创建安全组并加入服务主体再将安全组作为Contributor添加到目标工作区。托管身份Azure 托管部署适用于 Azure VM、AKS、App Service 等支持托管身份的 Azure 计算环境。托管身份同样需要被添加为工作区 Contributor且需启用上述租户设置该设置虽名为服务主体但同样适用于托管身份。# 系统分配的托管身份无需额外配置 credential: authentication_method: managed_identity# 用户分配的托管身份需提供 client_id credential: authentication_method: managed_identity managed_identity_client_id: your-managed-identity-client-idAzure CLI本地开发与测试使用本地az login会话凭据登录用户现有的 Fabric 权限直接生效。运行摄取前先执行az login远程无浏览器服务器可用az login。credential: authentication_method: cliDefaultAzureCredential灵活自动检测按顺序尝试多种凭据来源环境变量、工作负载身份、托管身份、共享令牌缓存、Azure CLI、Azure PowerShell、Azure Developer CLI 等。credential: authentication_method: default可通过exclude_*参数排除特定凭据来源加快检测速度或避免混合环境下的意外认证credential: authentication_method: default exclude_cli_credential: true # 生产环境建议跳过 Azure CLI exclude_environment_credential: false exclude_managed_identity_credential: false六、完整 Recipe 参考以下为官方提供的完整 Recipe 模板fabric-data-factory_recipe.yml可直接复制使用# Fabric Data Factory source 示例 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} # 可选按名称模式过滤工作区 workspace_pattern: allow: - .* # 默认允许所有工作区 deny: [] # 可选按名称模式过滤管道 pipeline_pattern: allow: - .* # 默认允许所有管道 deny: [] # 功能开关 extract_pipelines: true include_lineage: true include_execution_history: true execution_history_days: 7 # 1-90 天 # 可选将 Fabric 连接名映射到平台实例以获得准确血缘 # platform_instance_map: # my-snowflake-connection: prod_snowflake # my-bigquery-connection: analytics_project # 可选本连接器的平台实例多租户场景 # platform_instance: my-fabric-tenant # 环境 env: PROD # 可选启用状态化摄取以清理过期实体 # stateful_ingestion: # enabled: true sink: type: datahub-rest config: server: http://localhost:8080 token: ${DATAHUB_GMS_TOKEN}全部配置项说明对照 config.py连接器支持的完整配置项如下配置项类型默认值说明credential对象AzureCredentialConfigAzure 认证配置支持服务主体、托管身份、Azure CLI 与自动检测workspace_patternAllowDenyPattern全部允许按名称正则过滤工作区如allow[prod-.*], deny[.*-test]pipeline_patternAllowDenyPattern全部允许按名称正则过滤数据管道作用于所有匹配的工作区extract_pipelinesbooltrue是否提取数据管道及其活动include_lineagebooltrue是否从活动输入/输出提取血缘基于连接类型映射到 DataHub 数据集include_execution_historybooltrue是否将管道与活动执行历史提取为 DataProcessInstanceexecution_history_daysint7提取执行历史的天数1-90仅当include_execution_history为真时生效platform_instance_mapdict[str, str]{}连接名 → DataHub 平台实例的映射用于准确解析到已有数据集platform_instancestr无本连接器的平台实例多租户场景下防止 URN 冲突继承自PlatformInstanceConfigMixinenvstr无环境标识继承自EnvConfigMixinapi_timeoutint30REST API 调用超时秒数1-300stateful_ingestion对象None状态化摄取与过期实体清理配置启用后自动移除 Fabric 中已不存在的实体七、限制与边界以下是官方文档明确列出的能力边界务必在方案设计前确认通用限制运行历史上限Fabric API 每个管道最多返回 100 条最近完成的运行。若execution_history_days覆盖的运行数超过此上限只返回最近 100 条。建议提高摄取频率以捕获更久远的历史不支持 Dataflow Gen2Dataflow Gen2 项含转换逻辑的独立工作区级项不会被提取不支持 CopyJob工作区级独立 CopyJob 项不会被提取只有管道内嵌的 Copy 活动产生血缘无触发器/计划元数据管道的触发器与计划不会被提取ExecutePipeline 不支持ExecutePipeline活动类型在 Fabric 中已被标记为 legacy不支持用于跨管道血缘。血缘限制血缘范围只有 Copy 和 InvokePipeline 活动产生数据集或管道血缘其他活动类型Lookup、Wait、ForEach、Script 等作为 DataJob 摄取但无数据集级血缘InvokePipeline 操作类型仅支持InvokeFabricPipeline操作类型用于跨管道血缘InvokeAdfPipeline、InvokeExternalPipeline不会被解析将被跳过基于查询的 Copy 源当 Copy 活动使用sqlReaderQuery或sqlReaderStoredProcedureName而非直接表引用时不提取血缘无列级血缘连接器仅提取数据集级血缘Copy 活动翻译器配置中的列到列映射不会被提取无 Notebook/SparkJobDefinition 血缘Notebook 与 SparkJobDefinition 活动作为 DataJob 摄取但其血缘不会被解析连接解析未映射的连接类型回退为使用连接类型字符串作为平台名可能与 DataHub 中已有平台名不匹配请使用platform_instance_map显式映射连接名。八、故障排查401/403 错误确认服务主体具有正确的 Fabric API 权限并已添加为工作区成员Contributor。同时检查 Fabric 管理员的开发者设置是否已启用见上文Fabric 管理员设置注意设置变更传播最长需 15 分钟空结果检查workspace_pattern和pipeline_pattern是否过滤掉了所有项血缘缺失确认include_lineage: true已设置且管道中的 Fabric 连接配置正确。同时对照上文血缘限制章节检查活动类型与场景是否受支持实体过期残留启用stateful_ingestion以自动移除 Fabric 中已不存在的实体。九、测试与验证仓库中提供了针对该连接器的完整集成测试test_fabric_data_factory_source.py。测试使用 mock 的客户端响应验证完整摄取管线覆盖的场景包括工作区与管道过滤prod-analytics允许 /dev-sandbox过滤etl-daily允许 /test-validation过滤Copy 活动血缘Snowflake → LakehouseInvokePipeline 跨管道血缘执行历史1 条管道运行 成功/失败 2 条活动运行连接类型映射Snowflake 连接、Lakehouse 连接分别解析到外部数据集与 OneLake URN。测试通过Pipeline.create构造fabric-data-factory源配置并输出 golden 文件与预期元数据事件比对是理解本连接器实际行为的最佳参考实现。若你在配置后遇到血缘缺失或实体异常可以先对照该测试中的活动结构与连接数据结构排查。总结fabric-data-factory连接器把 Microsoft Fabric Data Factory 的编排元数据工作区、管道、活动、运行记录完整映射进 DataHub 的 Container / DataFlow / DataJob / DataProcessInstance 实体体系并以 Copy 与 InvokePipeline 活动为支点打通数据集级与跨管道血缘。理解其两步连接解析机制FABRIC_CONNECTION_PLATFORM_MAP自动映射 platform_instance_map手动映射、100 条运行上限、多租户 URN 隔离规则以及认证与权限前置条件是把它稳定运行在生产环境的关键。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表