ARTICLE DETAIL

资讯详情

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

DataHub Python SDK 完全指南:客户端、构建器与 Emitter 全解析

DataHub Python SDK 完全指南:客户端、构建器与 Emitter 全解析 DataHub Python SDK 完全指南客户端、构建器与 Emitter 全解析【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahubDataHub Python SDK 是 DataHub 元数据平台的官方 Python 客户端体系覆盖元数据写入Emitter、实体管理EntityClient、血缘建模LineageClient、搜索与解析SearchClient / ResolverClient以及事件构建MCE / MCP Builder等全链路能力。本指南以仓库中docs-website/sphinx/index.rst所定义的 SDK 文档骨架为主线结合metadata-ingestion/src/datahub下的真实源码与可运行示例帮助你快速掌握这套 SDK 的模块划分、核心 API 用法与底层实现原理并能直接在你的数据平台上落地数据集管理、血缘追踪和元数据上报。从 Sphinx 文档索引看 SDK 的模块体系docs-website/sphinx/index.rst是 DataHub Python SDK 文档的 Sphinx 主入口master document。它本身并不直接包含 API 讲解而是通过toctree指令把整个 SDK 的参考文档组织成六个区块这六块也正是 SDK 的能力地图toctree 区块覆盖模块对应源码目录apidocs/sdk-v2/*新一代 SDK 客户端主客户端、实体、血缘、搜索、解析、实体类metadata-ingestion/src/datahub/sdk/apidocs/builder/*MCE / MCP 事件构建器metadata-ingestion/src/datahub/emitter/mce_builder.py、mcp.py、mcp_builder.pyapidocs/clients/*Graph 客户端、Kafka Emitter、Rest Emittermetadata-ingestion/src/datahub/ingestion/graph/client.py、emitter/rest_emitter.py、emitter/kafka_emitter.pyapidocs/models元数据模型Aspect 类datahub.metadata.schema_classesapidocs/urnsURN 类型体系datahub.metadata.urns每个.rst文件如docs-website/sphinx/apidocs/sdk-v2/main-client.rst都通过 Sphinx 的automodule指令自动从 docstring 生成 API 文档例如.. automodule:: datahub.sdk.main_client :member-order: alphabetical这意味着docs-website/sphinx/conf.py中的 Sphinx 配置直接决定了文档产出质量。从 conf.py 可以看到该工程启用了sphinx.ext.autodoc、sphinx.ext.napoleon与sphinx_autodoc_typehints三个核心扩展napoleon_use_param True、autodoc_typehints description因此源码中的 Google 风格 docstring 与类型标注会被自动渲染为参数说明。文档网站侧通过 convert_sphinx_to_docusaurus.py 将 RST 产物转换后发布详见 docs-website/README.md 中的yarn run _generate-python-sdk生成流程。主客户端 DataHubClient三种初始化方式与子客户端SDK v2 的入口是datahub.sdk.main_client.DataHubClient见 main_client.py。官方文档《SDK v2 主客户端参考》将其定义为 a client for interacting with DataHub其 docstring 明确声明了三种初始化方式from datahub.sdk import DataHubClient # 方式 1直接指定 server URL 与可选 token client DataHubClient(serverhttp://localhost:8080, tokenyour_token) # 方式 2传入完整连接配置对象 from datahub.ingestion.graph.config import DatahubClientConfig client DataHubClient(configDatahubClientConfig(serverhttp://localhost:8080, token...)) # 方式 3包装已有的旧版DataHubGraph 实例 from datahub.ingestion.graph.client import DataHubGraph client DataHubClient(graphexisting_graph)源码中三种参数是互斥的同时传入server与config、server与graph都会抛出SdkUsageError三者皆空同样报错。此外DataHubClient.from_env()提供了环境变量/配置文件初始化路径——它优先读取DATAHUB_GMS_URL与DATAHUB_GMS_TOKEN环境变量否则回退到~/.datahubenv文件并支持client_mode默认ClientMode.SDK与datahub_component两个参数main_client.py。DataHubClient采用门面Facade式设计通过只读属性向调用方暴露一组分工明确的子客户端main_client.py属性类型职责client.entitiesEntityClient实体获取与增删改管理client.searchSearchClient按过滤条件检索实体 URNclient.lineageLineageClient血缘的写入、读取与 SQL 推断client.resolveResolverClient按名称/邮箱等解析出 URNclient.assertionsAssertionsClient断言管理云版优先缺失时回退实体客户端 EntityClient 与实体类EntityClient 是 client for retrieving and managing DataHub entities其构造函数为私有只能通过DataHubClient.entities获取实例。它对dataset、container、dashboard、chart、datajob、dataflow、mlmodel、mlmodelgroup、semantic_model、metric、document、glossary_term、glossary_node、tag等实体提供了类型化的管理能力。与实体客户端配套的是apidocs/sdk-v2/entities.rst中列举的 SDK 实体类。这些实体类位于metadata-ingestion/src/datahub/sdk/目录如dataset.py、container.py、dashboard.py、datajob.py、mlmodel.py、semantic_model.py等为每类元数据对象封装了属性读写与 URN 构造。例如用类型化 URN 直接构造数据集from datahub.metadata.urns import DatasetUrn urn DatasetUrn(platformsnowflake, namesales_raw) print(urn) # urn:li:dataset:(urn:li:dataPlatform:snowflake,sales_raw,PROD)血缘客户端 LineageClient写入、推断与读取docs-website/sphinx/apidocs/sdk-v2/lineage-client.rst明确指出 LineageClient 负责 adding and retrieving lineage information from DataHub并推荐阅读官方血缘教程作为更高层的入门材料。该教程lineage.md与本仓库示例文件共同构成了完整的血缘实践指南。建立实体血缘add_lineage()方法lineage_client.py接受upstream与downstream两个 URN支持数据集、数据任务、看板、图表之间的组合示例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)建立列级血缘fuzzy / strict / 自定义映射column_lineage参数控制列级血缘的匹配方式见血缘教程取值说明False关闭列级血缘默认True/auto_fuzzy按列名相似度自动匹配如user_id↔userId、customer_id↔CustomerIdauto_strict要求上下游列名完全一致字典自定义映射{下游列名: [上游列名列表]}# 自定义映射{ downstream_column - [upstream_columns] } client.lineage.add_lineage( upstreamDatasetUrn(platformsnowflake, namesales_raw), downstreamDatasetUrn(platformsnowflake, namesales_cleaned), column_lineage{ id: [id], region: [region, region_id], total_revenue: [revenue], }, )源码层面add_lineage内部根据上下游实体类型分发到_add_dataset_lineage、_add_dashboard_lineage、_add_datajob_lineage、_add_chart_lineage等私有方法lineage_client.py_get_strict_column_lineage与_get_fuzzy_column_lineagelineage_client.py分别实现精确匹配与基于名称规范化的模糊匹配。血缘教程还给出了一张支持的血缘组合表上游实体下游实体DatasetDatasetDatasetDataJobDataJobDataJobDataJobDatasetDatasetDashboardChartDashboardDashboardDashboardDatasetChart需要特别注意的是列级血缘与带transformation_text的查询节点仅支持 Dataset → Dataset组合。从 SQL 推断血缘infer_lineage_from_sql()lineage_client.py会解析 SQL 语句自动识别上游与下游数据集并尽可能生成列级血缘同时创建一个承载 SQL 变换逻辑的查询节点。教程 lineage.md 与示例 lineage_dataset_from_sql.py 给出了完整用法from 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 会被假定位于默认 db/schema即 prod_db.public.sales_summary client.lineage.infer_lineage_from_sql( query_textsql_query, platformsnowflake, default_dbprod_db, default_schemapublic, )更底层的 SQL 解析原理可参考DataHub SQL Parser 文档。读取血缘get_lineage()lineage_client.py按方向读取血缘默认只取 1 跳hop的直接上下游传入max_hops可跨多跳遍历当max_hops大于 2 时会遍历完整血缘图并按count限制结果。通过source_column或SchemaFieldUrn作为source_urn可获取列级血缘路径。示例 lineage_get_basic.pyfrom 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)返回值是一组LineageResult对象包含urn、type、hops、direction、platform、name字段列级查询额外带pathsLineagePath列表每个路径描述从源列到目标列的 URN 链。get_lineage同样支持按平台、类型、域、环境等条件过滤lineage_get_with_filter.py过滤器细节见下文搜索客户端的 Filter DSL。搜索客户端 SearchClient 与过滤器 DSLSearchClient 提供 searching and retrieving metadata from DataHub 的能力其核心方法是get_urns()search_client.py返回满足条件的实体 URN 迭代器。它底层委托给 Graph 客户端的get_urns_by_filter当前尚未实现分页参数源码中留有 TODO 注释。过滤条件由 search_filters.py 中的FilterDsl别名F构造并经compile_filters()search_client.py编译后传给后端。官方搜索教程 search_client.md 中详细说明了基于过滤器的搜索filter-based search。典型的过滤器组合方式如下from datahub.sdk.search_filters import FilterDsl as F filters F.and_( F.entity_type(dataset), F.or_( F.custom_filter(platform, EQUAL, [snowflake]), F.custom_filter(platform, EQUAL, [bigquery]), ), ) urns client.search.get_urns(filterfilters)ResolverClientresolver_client.py正是构建在这一搜索能力之上的domain(name...)、user(name...)/user(email...)、term(name...)等方法先按名称或邮箱等条件搜索实体再将其转换为对应的DomainUrn、CorpUserUrn、GlossaryTermUrn。找不到时抛出ItemNotFoundError匹配到多个则抛出MultipleItemsFoundErrorresolver_client.py。事件构建器MCE Builder 与 MCP Builderapidocs/builder区块包含两类构建器用于构造 DataHub 的元数据变更事件。MCE Builder构造 MetadataChangeEventmce_builder.py 提供了一系列工具函数make it easier to construct MetadataChangeEvents。其中最常用的 URN 构建函数包括函数作用make_data_platform_urn(platform)生成数据平台 URN如urn:li:dataPlatform:snowflakemake_dataset_urn(platform, name, envDEFAULT_ENV)生成数据集 URN默认环境为PRODmake_dataset_urn_with_platform_instance(...)带平台实例的数据集 URNmake_schema_field_urn(parent_urn, field_path)生成 schema 字段列URNmake_container_urn(guid)生成容器 URN基于 GUIDmake_user_urn(username)/make_group_urn(groupname)用户与用户组 URNmake_owner_urn(owner, owner_type)按所有权类型生成所有者 URNmake_term_urn(term)/make_tag_urn(tag)术语与标签 URNmake_ts_millis(ts)/parse_ts_millis(ts)时间戳与毫秒时间戳互转datahub_guid(obj)基于字典内容生成确定性 GUID容器 URN 的基础MCP Builder构造 MetadataChangeProposalmcp.py 中的MetadataChangeProposalWrapper是 MCPMetadataChangeProposal事件的便捷封装mcp_builder.py 则提供了面向平台/容器等复杂场景的构建工具。与 MCE 相比MCP 是目前推荐的增量式元数据变更格式可通过 Rest/Kafka emitter 直接发送。MetadataPatchProposal位于 mcp_patch_builder.py则支持基于 patch 的部分更新。Emitter 与 Graph Client三条元数据写入通道apidocs/clients区块对应三类底层连接组件是 SDK 与 DataHub 服务端交互的通道。Rest EmitterDataHubRestEmitter can be used to push metadata to DataHub通过 HTTP 调用 GMS 的 REST 端点。核心方法emit()统一入口自动识别 MCE / MCP 并路由rest_emitter.pyemit_mce(mce)/emit_mcp(mcp)分别发送 MCE 与 MCPrest_emitter.pyemit_mcps(...)批量发送并自动分块chunk大负载场景由_DEFAULT_EMIT_MODE决定走 OpenAPI 还是 RestLi 端点rest_emitter.pyemit_usage(usageStats)上报使用量统计rest_emitter.pytest_connection()连通性与认证自检rest_emitter.pyKafka EmitterDataHubKafkaEmitter 面向高吞吐场景将 MCE / MCP 写入 Kafka topic 而非直接调用 HTTP适合批量、异步的元数据管道。它提供emit()、emit_mce_async()与emit_mcp_async()kafka_emitter.py构造时需传入KafkaEmitterConfigbootstrap 地址、topic、序列化配置等。Graph ClientDataHubGraph 继承自DatahubRestEmitter因此同时具备发送能力与查询能力——这正是apidocs/clients/graph-client.rst所描述的 extends the Rest emitter with additional functionality。其典型查询方法包括get_aspect()client.py、get_aspect_v2()client.py、get_aspects_for_entity()、get_aspect_counts()以及get_urns_by_filter()等覆盖了 aspect 读取、实体 URN 检索等元数据服务能力。在DataHubClient内部无论是 server/token 直连还是from_env()最终都会包装为一个DataHubGraph实例作为底层_graph。模型与 URN 体系Modelsapidocs/models.rstdatahub.metadata.schema_classes是 DataHub 全部元数据模型Aspect 类的自动生成命名空间construct_with_defaults、RECORD_SCHEMA、ASPECT_NAME、ASPECT_INFO等内部实现细节被 Sphinx 排除在文档之外便于聚焦面向使用者的字段与构造器。URNsapidocs/urns.rstdatahub.metadata.urns提供DatasetUrn、DataFlowUrn、DataJobUrn、ChartUrn、DashboardUrn、ContainerUrn、MlModelUrn、MlModelGroupUrn、MetricUrn、SemanticModelUrn、DocumentUrn、GlossaryNodeUrn、GlossaryTermUrn、CorpUserUrn、DomainUrn等类型化 URN 类全部导入于 entity_client.py。这些类封装了 URN 的构造、解析与校验LI_DOMAIN、URN_PREFIX、get_entity_id、url_encode等底层细节同样被文档排除。典型落地流程组合使用各模块把上述模块串起来一个典型的注册元数据 建立血缘 验证检索工作流大致如下from datahub.sdk import DataHubClient from datahub.metadata.urns import DatasetUrn client DataHubClient(serverhttp://localhost:8080, tokenyour_token) # 1) 建立数据集血缘含列级 fuzzy 匹配 client.lineage.add_lineage( upstreamDatasetUrn(platformsnowflake, namesales_raw), downstreamDatasetUrn(platformsnowflake, namesales_cleaned), column_lineageTrue, ) # 2) 读取下游血缘 print(client.lineage.get_lineage( source_urnDatasetUrn(platformsnowflake, namesales_cleaned), directiondownstream, )) # 3) 用搜索客户端验证实体是否已入库 from datahub.sdk.search_filters import FilterDsl as F urns list(client.search.get_urns( filterF.and_( F.entity_type(dataset), F.custom_filter(platform, EQUAL, [snowflake]), ) )) print(urns)总结与延伸阅读DataHub Python SDK 的文档体系与源码结构一一对应sdk-v2区块覆盖新一代DataHubClient门面及其五个子客户端builder区块提供 MCE/MCP 事件构建能力clients区块提供 Rest/Kafka 两条写入通道与DataHubGraph查询能力models与urns则承载底层数据模型。对于需要深入某一能力的读者建议继续阅读血缘使用教程 与搜索 SDK 教程可运行示例lineage_dataset_add.py、lineage_dataset_column.py、lineage_get_with_hops.py、lineage_dataset_from_sql.pySDK 源码main_client.py、lineage_client.py、search_client.py、rest_emitter.py、graph/client.py文档生成基础设施conf.py 与 convert_sphinx_to_docusaurus.py【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表