
DataHub 集成 Microsoft Fabric OneLake容器建模、视图血缘与查询用量提取实战指南【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub本文以 DataHub 仓库中的fabric-onelake摄取源为核心系统讲解如何将 Microsoft Fabric OneLake 的工作区Workspace、湖仓Lakehouse、仓库Warehouse、Schema 与表/视图映射为 DataHub 实体并利用 SQL Analytics Endpoint 提取列级 Schema、从视图定义解析视图→表血缘以及从queryinsights视图采集查询用量统计与操作事件。读完本文你将掌握 OneLake 连接器的完整配置参数、认证与权限准备、运行方式及其底层实现原理。功能概览DataHub 对 Microsoft Fabric OneLake 的集成覆盖了工作区、湖仓、仓库三类容器以及带 Schema 元数据的表数据集Dataset与带视图定义、且能通过 SQL 解析产出视图→表血缘的视图数据集同时从 SQL Analytics Endpoint 的queryinsights视图中提取查询用量统计并支持有状态的删除检测stateful deletion detection。从源码实现看该连接器的入口为 source.py其中声明了FabricOneLakeSource的能力清单source.py#L136-L155支持容器Container提取默认启用Schema 元数据提取默认启用平台实例Platform Instance默认启用细粒度血缘在启用extract_views时通过 SQL 解析视图定义获得用量统计在启用usage.include_usage_statistics时读取queryinsights.exec_requests_history30 天保留期并通过 SQL 解析推导列级用量操作事件捕获可通过usage.include_operational_stats可选开启。概念映射与层级结构连接器将 Fabric 对象映射为 DataHub 实体的方式如下Microsoft FabricDataHub 实体说明Workspace工作区Container子类型Fabric Workspace顶层组织单元Lakehouse湖仓Container子类型Fabric Lakehouse包含 Schema 与表Warehouse仓库Container子类型Fabric Warehouse包含 Schema 与表SchemaContainer子类型Fabric Schema湖仓/仓库内部的逻辑分组Table表DatasetSchema 内的表View视图Dataset子类型View湖仓与仓库视图通过 SQL 解析从视图定义提取血缘在源码中这些子类型来自 source.py 对DatasetContainerSubTypes.FABRIC_LAKEHOUSE/FABRIC_WAREHOUSE/FABRIC_SCHEMA以及DatasetSubTypes.TABLE/DatasetSubTypes.VIEW的使用容器通过LakehouseKey、WarehouseKey、LakehouseSchemaKey、WarehouseSchemaKey等容器键ContainerKey建立父子层级实现parent_key()遍历。层级结构Platform (fabric-onelake) └── Workspace (Container) ├── Lakehouse (Container) │ └── Schema (Container) │ └── Table/View (Dataset) └── Warehouse (Container) └── Schema (Container) └── Table/View (Dataset)平台实例即租户Fabric REST API 不暴露租户级端点。为了在 DataHub 中表达租户级组织请在配置中把platform_instance字段设置为你的租户标识符例如contoso-tenant。该值会进入所有容器与数据集的 URN从而把所有工作区有效地归并到指定的平台实例/租户之下。这一设计同时写在了配置类的文档字符串中config.py#L129-L136所有 API 操作都在工作区级别进行租户组织需要借助platform_instance来表达。前置条件开始摄取前请确保满足网络连通性、有效的认证凭据以及元数据 API 所需的读权限。认证方式连接器支持多种 Azure 认证方法通过credential.authentication_method指定方法适用场景配置Service Principal服务主体生产环境authentication_method: service_principalManaged Identity托管标识Azure 托管部署VM、AKS、App Service 等authentication_method: managed_identityAzure CLI本地开发authentication_method: cli先执行az loginDefaultAzureCredential灵活多变的环境authentication_method: default认证底层由 common/auth.py 的FabricAuthHelper负责配置模型复用AzureCredentialConfigconfig.py#L139-L146支持服务主体、托管标识、Azure CLI 及自动探测DefaultAzureCredential。所需权限连接器需要对 Fabric 工作区及其内容拥有只读访问权。被认证的身份服务主体、托管标识或用户必须具备工作区级权限Workspace.Read.All或Workspace.ReadWrite.AllMicrosoft Entra 委派作用域在要摄取的工作区中拥有Viewer角色或更高角色。API 权限Entra API 权限Workspace.Read.All委派——列出并读取工作区元数据或Workspace.ReadWrite.All委派——读写访问。Token 受众Token Audiences连接器按操作使用两种不同的 token 受众Fabric REST APIhttps://api.fabric.microsoft.com使用 Power BI API 作用域https://analysis.windows.net/powerbi/api/.default用于列出工作区、湖仓、仓库及基础表元数据OneLake Delta Table APIhttps://onelake.table.fabric.microsoft.com使用 Storage 受众https://storage.azure.com/.default用于访问启用 Schema 的湖仓中的 Schema 与表。连接器会自动处理两种 token 受众启用 Schema 的湖仓走 OneLake Delta Table APIStorage 受众 token未启用 Schema 的湖仓走标准 Fabric REST API。无需额外配置。OneLake 数据访问权限对启用 Schema 的湖仓若该湖仓启用了 OneLake 安全控制请确保身份对湖仓项拥有Read或ReadWrite权限。这些权限独立于工作区角色需在 Fabric 门户中湖仓的安全设置里管理。SQL Analytics Endpoint 设置通过 SQL Analytics Endpoint 做 Schema 提取需要在系统上安装 ODBC 驱动。1. 安装 ODBC Driver ManagerUnixODBCUbuntu/Debiansudo apt-get install -y unixodbc unixodbc-devRHEL/CentOS 7/8sudo yum install -y unixODBC unixODBC-develFedora / RHEL 9sudo dnf install -y unixODBC unixODBC-develmacOSbrew install unixodbc2. 安装 Microsoft ODBC Driver 18 for SQL Server连接 Fabric SQL Analytics Endpoint 所必需各 Linux/macOS 发行版安装命令详见 fabric-onelake_pre.md。3. 验证安装odbcinst -q -d列表中应出现ODBC Driver 18 for SQL Server。4. 权限你的 Azure 身份需具备查询 SQL Analytics Endpoint 的权限与使用 SQL 工具访问该端点的权限一致。5. Python 依赖fabric-onelakeextra 包含sqlalchemy与pyodbcpip install acryl-datahub[fabric-onelake]若遇到libodbc.so.2: cannot open shared object file错误请先确认已安装 ODBC Driver Manager上述第 1 步。视图提取的VIEW DEFINITION权限视图提取复用 SQL Analytics Endpoint 连接同一套 ODBC 驱动且只有在sql_endpoint.enabled为true时才会尝试提取视图。读取视图定义视图→表血缘的前提需要 SQL Analytics Endpoint 上的VIEW DEFINITION权限——表摄取所用的工作区Viewer角色并不够它只授予db_datareader这会使INFORMATION_SCHEMA.VIEWS.VIEW_DEFINITION返回NULL。该权限没有工作区级开关必须二选一按湖仓/仓库单独授予VIEW DEFINITION最小权限推荐身份保持工作区 ViewerGRANT VIEW DEFINITION ON DATABASE::lakehouse_or_warehouse_name TO [service_principal_name];赋予更高的工作区角色Contributor、Member 或 Admin。若两种方案在你的环境中都不可行可将extract_views: false以跳过视图摄取。此时若仍以 Viewer 级别摄取视图视图会出现但没有定义血缘将缺失。查询用量统计的权限与保留期用量提取通过 SQL Analytics Endpoint 读取queryinsights.exec_requests_history复用上述 ODBC 设置因此sql_endpoint.enabled必须为true——配置校验器会拒绝usage.include_usage_statisticstrue而 SQL Endpoint 未启用的组合该校验逻辑见 config.py#L266-L299。所需角色queryinsights的可见性按工作区作用域划分。摄取身份服务主体、托管标识或用户需要在每个目标工作区拥有Contributor 或更高角色。用于表摄取的工作区 Viewer 角色不够queryinsights要求 Premium 容量工作区上具备 contributor or higher 权限且完整查询文本SQL 解析与列级用量所必需只对 Admin、Member 和 Contributor 角色开放。保留期与延迟Fabric 仅保留queryinsights30 天——更早的历史无法回填请据此配置usage.start_time。新执行的查询最多需要 15 分钟才会出现并发较高时延迟会增大。系统查询及用户上下文之外的查询不会出现在结果中。运行摄取基础 Recipe参考 fabric-onelake_recipe.yml 作为模板source: type: fabric-onelake 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: Platform instance (use as tenant identifier) # This groups all workspaces under a tenant-level container # platform_instance: contoso-tenant # Optional: Environment # env: PROD # Optional: Filter workspaces by name pattern # workspace_pattern: # allow: # - .* # Allow all workspaces by default # deny: [] # Optional: Filter lakehouses by name pattern # lakehouse_pattern: # allow: # - .* # Allow all lakehouses by default # deny: [] # Optional: Filter warehouses by name pattern # warehouse_pattern: # allow: # - .* # Allow all warehouses by default # deny: [] # Optional: Filter tables by name pattern # Format: schema.table or just table for default schema # table_pattern: # allow: # - .* # Allow all tables by default # deny: [] # Optional: Filter views by name pattern # Format: schema.view or just view for default schema # view_pattern: # allow: # - .* # Allow all views by default # deny: [] # Feature flags extract_lakehouses: true extract_warehouses: true extract_schemas: true # Set to false to skip schema containers extract_views: true # Requires sql_endpoint.enabled # Optional: API timeout (seconds) # api_timeout: 30 # Optional: Stateful ingestion for stale entity removal # stateful_ingestion: # enabled: true # remove_stale_metadata: true sink: type: datahub-rest config: server: http://localhost:8080执行摄取命令datahub ingest -c fabric-onelake_recipe.yml高级配置source: type: fabric-onelake config: credential: authentication_method: service_principal client_id: ${AZURE_CLIENT_ID} client_secret: ${AZURE_CLIENT_SECRET} tenant_id: ${AZURE_TENANT_ID} # Platform instance (represents tenant) platform_instance: contoso-tenant # Environment env: PROD # Filtering workspace_pattern: allow: - prod-.* - shared-.* deny: - .*-test - .*-dev lakehouse_pattern: allow: - .* deny: - .*-backup warehouse_pattern: allow: - .* deny: [] table_pattern: allow: - .* deny: - .*_temp - .*_backup view_pattern: allow: - .* deny: - .*_internal # Feature flags extract_lakehouses: true extract_warehouses: true extract_schemas: true # Set to false to skip schema containers extract_views: true # Requires sql_endpoint.enabled # API timeout (seconds) api_timeout: 30 # Stateful ingestion (optional) stateful_ingestion: enabled: true remove_stale_metadata: true sink: type: datahub-rest config: server: http://localhost:8080使用托管标识Managed Identitysource: type: fabric-onelake config: credential: authentication_method: managed_identity # For user-assigned managed identity, specify client_id # client_id: ${MANAGED_IDENTITY_CLIENT_ID} platform_instance: contoso-tenant env: PROD sink: type: datahub-rest config: server: http://localhost:8080使用 Azure CLI本地开发source: type: fabric-onelake config: credential: authentication_method: cli # Run az login first platform_instance: contoso-tenant env: DEV sink: type: datahub-rest config: server: http://localhost:8080配置参数详解以下参数来自 config.py 中的FabricOneLakeSourceConfig及子配置模型均可用 pydantic 校验过滤模式AllowDenyPattern均为正则参数说明workspace_pattern按名称过滤工作区例allow[prod-.*], deny[.*-test]lakehouse_pattern按名称过滤湖仓作用于所有匹配workspace_pattern的工作区warehouse_pattern按名称过滤仓库作用于所有匹配workspace_pattern的工作区schema_pattern按名称过滤 Schema被拒绝的 Schema 整体跳过含其全部表与视图table_pattern按名称过滤表格式schema.table或默认 Schema 下仅用tableview_pattern按名称过滤视图格式schema.view_name或默认 Schema 下仅用view_name功能开关参数默认值说明extract_lakehousestrue是否提取湖仓及其表extract_warehousestrue是否提取仓库及其表extract_viewstrue是否提取视图及定义依赖sql_endpoint视图通过 SQL Analytics Endpoint 上的INFORMATION_SCHEMA.VIEWS发现extract_schemastrue是否提取 Schema 容器为false时表直接挂在湖仓/仓库容器下api_timeout30REST API 调用超时秒范围 1–300Schema 提取extract_schema参数默认值说明enabledtrue启用 Schema 提取methodsql_analytics_endpoint提取方式当前仅支持sql_analytics_endpointSQL Analytics Endpointsql_endpoint参数默认值说明enabledtrue启用 SQL Analytics Endpoint 连接odbc_driverODBC Driver 18 for SQL ServerSQL Server 连接使用的 ODBC 驱动名encryptyes连接加密合法值yes/mandatory启用加密ODBC Driver 18.0 默认、no/optional禁用、strictODBC Driver 18.0、仅 TDS 8.0 协议始终校验服务器证书trust_server_certificateno是否免校验信任服务器证书仅在证书校验失败时设为yesencryptstrict时该设置被忽略、始终校验query_timeout30SQL 查询超时秒范围 1–300用量统计usage继承BaseUsageConfig参数默认值说明include_usage_statisticstrue用量提取总开关为false时不从queryinsights发出任何datasetUsageStatistics或operation方面skip_failed_queriestrue为true时 SQL 过滤掉status ! Succeeded的行已取消/失败的查询在源头跳过include_queriestrue为每个去重后的queryinsights行发出Query实体使 SQL 成为 DataHub 中可检索的一等资产Queries 标签页与独立 Query 页面可见需include_usage_statisticsTrue此外usage块还支持全部标准BaseUsageConfig字段bucket_duration、start_time、end_time、top_n_queries、format_sql_queries、include_top_n_queries、include_operational_stats、user_email_pattern等。启用状态化摄取后用量时间窗口只有在运行成功后才会写入检查点checkpoint部分失败或失败的运行不会静默跳过下一个窗口。依赖校验validate_sql_endpoint_dependenciesconfig.py#L266-L299规定当extract_viewsTrue、extract_schema.methodsql_analytics_endpoint、或usage.include_usage_statisticsTrue三者任一成立时sql_endpoint必须配置且enabledTrue否则配置校验直接抛错——这三类功能都依赖 SQL Analytics Endpoint 查询。Schema 提取原理Schema 提取列元数据通过 SQL Analytics Endpoint 完成从湖仓与仓库的表/视图中提取列名、数据类型、可空性与序号位置。source: type: fabric-onelake config: credential: authentication_method: service_principal client_id: ${AZURE_CLIENT_ID} client_secret: ${AZURE_CLIENT_SECRET} tenant_id: ${AZURE_TENANT_ID} # Schema extraction configuration extract_schema: enabled: true # Enable schema extraction (default: true) method: sql_analytics_endpoint # Currently only this method is supported # SQL Analytics Endpoint configuration sql_endpoint: enabled: true # Enable SQL endpoint connection (default: true) # Optional: ODBC connection options # odbc_driver: ODBC Driver 18 for SQL Server # Default: ODBC Driver 18 for SQL Server # encrypt: yes # Enable encryption (default: yes) # trust_server_certificate: no # Trust server certificate (default: no) query_timeout: 30 # Timeout for SQL queries in seconds (default: 30)工作流程端点发现SQL Analytics Endpoint URL 自动从 Fabric API 为每个湖仓/仓库获取。端点格式为unique-identifier.datawarehouse.fabric.microsoft.com无法仅凭 workspace_id 构造。认证复用 REST API 访问配置的同一套 Azure 凭据注入 Azure AD token。连接用发现的端点 URL 通过 ODBC 连接 SQL Analytics Endpoint。查询查询INFORMATION_SCHEMA.COLUMNS提取列元数据Schema 提取所必需。类型映射SQL Server 数据类型由 Dataset SDK 通过resolve_sql_type()自动映射为 DataHub 类型见 source.py#L685-L699 中原始类型字符串传参的注释说明。重要说明端点 URL 必须从 API 获取若获取失败该项的 Schema 提取会失败。与旧版 Power BI Premium 端点不同Fabric SQL Analytics Endpoint 不支持 fallback 连接串。禁用 Schema 提取表将不带列元数据摄取source: type: fabric-onelake config: extract_schema: enabled: false视图提取与视图→表血缘湖仓与仓库中的视图会以View子类型的 DataHubDataset实体摄取每个视图数据集包含列级 Schema 元数据与表列一同来自INFORMATION_SCHEMA.COLUMNS原始视图定义CREATE VIEWSQL取自INFORMATION_SCHEMA.VIEWS由 SQL 解析聚合器SqlParsingAggregator从视图定义解析出的上游表血缘。source: type: fabric-onelake config: # View extraction is enabled by default. Set to false to skip views. extract_views: true # Filter views by name pattern. Format: schema.view or just view for default schema. view_pattern: allow: - .* deny: - .*_internal # View extraction requires the SQL Analytics Endpoint (enabled by default). sql_endpoint: enabled: true工作流程发现连接器在 SQL Analytics Endpoint 上查询INFORMATION_SCHEMA.VIEWS列出视图并捕获其定义。过滤每个视图按schema.view_name形式与view_pattern匹配。Schema列元数据复用同一份INFORMATION_SCHEMA.COLUMNS查询结果不为每个视图额外发起查询。血缘视图定义被送入 SQL 解析聚合器推导出视图 → 上游表血缘视图 URN 与上游表 URN 在同一工作区与项item内解析。从源码看source.py#L980-L991视图数据集以parse_view_lineageFalse发出随后通过aggregator.add_view_definition(...)将定义连同default_dbworkspace.id.item_id与规范化后的default_schema交给聚合器统一解析血缘在摄取结束时统一排空drain以便跨项的视图→表引用都能解析成功。查询用量统计连接器读取每个湖仓与仓库 SQL Analytics Endpoint 上的queryinsights.exec_requests_history视图提取查询用量。每条捕获的查询经 SQL 解析聚合器解析后输出为datasetUsageStatistics方面——查询计数、去重用户数、Top 用户、Top 字段以及启用时Top SQL 查询按配置的时间桶bucket聚合operation方面——启用usage.include_operational_stats时按查询输出操作事件insert、update、delete 等。source: type: fabric-onelake config: # Usage extraction is enabled by default. Set to false to skip query usage. usage: include_usage_statistics: true # When true, the SQL filter excludes rows where status ! Succeeded # (canceled / failed queries are skipped at the source). skip_failed_queries: true # Optional: emit per-query operation aspects in addition to aggregated # datasetUsageStatistics. Defaults to true (inherited from BaseUsageConfig). include_operational_stats: true # Optional: include top SQL queries in the usage payload. include_top_n_queries: true top_n_queries: 10 # Optional: window the connector queries from queryinsights. Defaults to # the standard BaseUsageConfig last bucket window. Fabric retains # queryinsights for 30 days. bucket_duration: DAY # start_time: 2026-04-01T00:00:00Z # end_time: 2026-05-01T00:00:00Z # Usage extraction depends on the SQL Analytics Endpoint. extract_schema: enabled: true sql_endpoint: enabled: true实现要点见 usage.pyFabricUsageExtractor将queryinsights的行流式送入共享的SqlParsingAggregator聚合器用 sqlglot 解析每条查询并产出datasetUsageStatistics/operation方面时间窗口通过RedundantUsageRunSkipHandler解析启用状态化摄取时跳过已被此前成功运行完整覆盖的窗口。注意queryinsights仅保留 30 天超出usage.start_time配置的历史无法回填新查询最多 15 分钟才可见。若某个湖仓/仓库的端点不可达该项目的用量会被跳过而不会导致整个运行失败。Schema 启用与未启用的湖仓连接器自动处理两类湖仓无需改动配置Schemas-Enabled 湖仓连接器先通过 OneLake Delta Table API 列出 Schema再列出每个 Schema 内的表。这需要 Storage 受众 tokenhttps://storage.azure.com/.default。Schemas-Disabled 湖仓连接器使用标准 Fabric REST API 的/tables端点列出全部表。没有显式 Schema 的表自动归入 DataHub 中的dboSchema该默认 Schema 常量定义于 constants.py 的FABRIC_SQL_DEFAULT_SCHEMA dbo适用于 Fabric SQL 层包括未启用 Schema 的湖仓与未限定的表引用。此路径使用 Power BI API 作用域 token。重要即使对未启用 Schema 的湖仓DataHub 中所有表的 URN 也都会带有 Schema——没有显式 Schema 的表统一归一化为dbo从而保证所有 Fabric 实体 URN 结构一致。状态化摄取与删除检测连接器支持状态化摄取以跟踪已摄取实体并移除过期元数据stateful_ingestion: enabled: true remove_stale_metadata: true启用后连接器将跟踪所有已摄取的工作区、湖仓、仓库、Schema 与表从 DataHub 移除在 Fabric 中已不存在的实体跨摄取运行维护状态。这一能力由StatefulStaleMetadataRemovalConfigconfig.py#L246-L254接入FabricOneLakeSource继承自StatefulIngestionSourceBase并在构造函数中按配置决定是否创建RedundantUsageRunSkipHandler将用量窗口已覆盖的判定与删除检测状态统一管理。限制元数据同步延迟SQL Analytics Endpoint 反映 Schema 变更可能存在延迟新列或 Schema 修改可能需要数分钟到数小时才出现。缺失的表部分表可能在 SQL 端点不可见原因包括不支持的数据类型、权限问题、超大数据库中表数量上限。优雅降级若某张表的 Schema 提取失败该表仍会不带列元数据照常摄取不导致摄取失败。视图提取依赖 SQL 端点视图只能通过 SQL Analytics Endpoint 发现若sql_endpoint.enabled为false或某湖仓/仓库端点不可达该项的视图不会被摄取。用量保留期Fabricqueryinsights仅保留 30 天查询历史无论usage.start_time如何配置更早的用量都无法回填。用量依赖 SQL 端点若sql_endpoint.enabled为false配置校验器会拒绝usage.include_usage_statisticstrue若某个湖仓/仓库端点不可达该项目用量被跳过但不会导致运行失败。故障排查若摄取失败请依次核对凭据、权限、连通性与作用域过滤再结合摄取日志中的源特定错误调整配置。常见问题方向libodbc.so.2: cannot open shared object file→ 安装 ODBC Driver Manager视图无定义/血缘缺失 → 检查VIEW DEFINITION权限或工作区角色无用量数据 → 确认身份具备工作区 Contributor 或更高角色且sql_endpoint.enabled: true配置校验报错 → 确认extract_views、extract_schema、usage与sql_endpoint的依赖关系满足三者任一开启都需要启用 SQL Endpoint。源码与测试验证想深入理解实现可继续阅读配置模型config.py摄取主流程与容器/数据集构建source.py用量/操作提取器usage.pySQL Analytics Endpoint 客户端schema_client.py平台标识与连接类型映射common/constants.py仓库中的集成测试与 golden 文件可以验证上述行为测试用例位于 test_fabric_onelake_source.py配套的 golden 文件覆盖带表的湖仓带视图的湖仓带视图的仓库含用量统计四类典型场景golden 目录下test_fabric_onelake_*_golden.json可用于对照实际输出验证容器层级、视图血缘与用量方面的结构是否符合预期。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考