news 2026/9/16 15:27:35

DataHub Semantic Models 实战:用 Python SDK 构建语义模型、逻辑数据集与指标血缘

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
DataHub Semantic Models 实战:用 Python SDK 构建语义模型、逻辑数据集与指标血缘

DataHub Semantic Models 实战:用 Python SDK 构建语义模型、逻辑数据集与指标血缘

【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

本文基于 DataHub 官方教程docs/api/tutorials/semantic-models.md,讲解如何使用 Python SDK v2(datahub.sdk)在 DataHub 中发出 Semantic Model(语义模型)、逻辑 Dataset(语义模型数据集)和 Metric(指标)实体:如何搭建Metric → Logical Dataset → Physical Dataset的完整血缘链、SDK 替你生成了哪些 aspect、如何校验输出的 MCP 形状,以及如何通过预检助手规避服务端版本不兼容问题。读完后,你可以直接复制示例代码向任意支持 semantic-model 元数据模型的服务端发出语义层元数据。

为什么要用 Semantic Models 与 Metrics

Semantic Models 和 Metrics 让你在物理数据之上描述一个逻辑层

  • 一个semanticModel实体将多个逻辑dataset(每个都是 "Semantic Model Dataset" 子类型)分组在一起;
  • 通过锚定在 schema field 上的semanticFieldAnnotation,暴露维度(dimensions)和度量(measures);
  • 它作为metric实体的底层模型(backing model)。

血缘流向为Metric → Logical Dataset → Physical Dataset(SemanticModel 是其成员的容器,不是血缘跳),从而为分析、治理和 AI 辅助探索提供一个稳定的、与数据源无关的表面。

这套建模方式最初存在于 Snowflake 连接器内部,被提升到高层 SDK 之后,任何生产者(连接器或直接使用 SDK 的用户)都可以发出同样的实体,而无需重新实现 aspect 装配逻辑。

本指南的目标

  • 构建一个包含两个逻辑数据集、schema 字段和一条 relationship 的SemanticModel
  • 发出由该模型支撑的metric实体,其中包含一个从另一个指标派生(derived)的指标;
  • 将每个发出的 MCP 序列化到文件,检查最终的 aspect 形状。

前置条件

本教程需要安装 DataHub SDK v2(datahub.sdk.*)。如果你是通过metadata-ingestion包运行示例,虚拟环境可通过./gradlew :metadata-ingestion:installDev配置好。

构建 Semantic Model、逻辑数据集与指标

下面的示例使用高层datahub.sdkbuilder 构建完整的血缘链Metric -> Logical Dataset -> Physical Dataset(SemanticModel 作为其 datasets 与 metrics 的容器),然后把每个发出的 MCP 写入 JSON 文件,以便检查 aspect 形状。完整可运行版本见 示例脚本,运行方式为python -m examples.library.semantic_model_create

"""Emit a semantic model with two logical datasets and two metrics. Run with: python -m examples.library.semantic_model_create """ import json from typing import Any, List from datahub.emitter.mce_builder import make_dataset_urn from datahub.emitter.mcp import MetadataChangeProposalWrapper from datahub.metadata.schema_classes import ( DialectClass, ERModelRelationshipCardinalityClass, SemanticFieldTypeClass, ) from datahub.metadata.urns import SemanticModelUrn from datahub.sdk import ( AiContextInput, DialectExpressionInput, Metric, SemanticFieldInput, SemanticModel, SemanticModelDataset, SemanticModelRelationshipInput, ) from datahub.sdk.entity import Entity def build_graph() -> tuple[SemanticModel, List[Entity]]: platform = "snowflake" model_urn = SemanticModelUrn(platform=platform, path="analytics", id="orders_model") orders_ds = SemanticModelDataset( platform=platform, name="analytics.orders_model.orders_ds", semantic_model=model_urn, alias="ORDERS", schema=[ SemanticFieldInput( field_path="order_id", type="int", semantic_type=SemanticFieldTypeClass.DIMENSION, is_part_of_key=True, ), # Foreign key the ORDERS -> CUSTOMERS relationship joins on. SemanticFieldInput( field_path="customer_id", type="int", semantic_type=SemanticFieldTypeClass.DIMENSION, ), SemanticFieldInput( field_path="order_ts", type="timestamp", semantic_type=SemanticFieldTypeClass.DIMENSION, is_time_dimension=True, ), SemanticFieldInput( field_path="amount", type="float", semantic_type=SemanticFieldTypeClass.MEASURE, expression=DialectExpressionInput( expression="SUM(amount)", dialect=DialectClass.SNOWFLAKE ), aggregation_function="SUM", ai_context=AiContextInput(synonyms=["revenue"]), ), ], upstreams=[make_dataset_urn(platform, "raw.orders")], ) customers_ds = SemanticModelDataset( platform=platform, name="analytics.orders_model.customers_ds", semantic_model=model_urn, alias="CUSTOMERS", schema=[ SemanticFieldInput( field_path="customer_id", type="int", semantic_type=SemanticFieldTypeClass.DIMENSION, is_part_of_key=True, ), SemanticFieldInput( field_path="customer_name", type="varchar", semantic_type=SemanticFieldTypeClass.DIMENSION, ), ], upstreams=[make_dataset_urn(platform, "raw.customers")], ) total_revenue = Metric( platform=platform, path="analytics", id="total_revenue", semantic_model=str(model_urn), name="Total Revenue", description="Sum of all order amounts.", expression=DialectExpressionInput( expression="SUM(ORDERS.amount)", dialect=DialectClass.SNOWFLAKE ), upstream_datasets=[orders_ds.urn], ai_context=AiContextInput(synonyms=["revenue"]), ) double_revenue = Metric( platform=platform, path="analytics", id="double_revenue", semantic_model=str(model_urn), name="Double Revenue", expression="2 * total_revenue", derived_from=[total_revenue.urn], ) model = SemanticModel( platform=platform, path="analytics", id="orders_model", name="Orders Model", description="A semantic model over the raw orders and customers tables.", datasets=[orders_ds, customers_ds], relationships=[ SemanticModelRelationshipInput( from_alias="ORDERS", from_columns=["customer_id"], to_alias="CUSTOMERS", to_columns=["customer_id"], name="orders_to_customers", cardinality=ERModelRelationshipCardinalityClass.N_ONE, ) ], ai_context=AiContextInput( synonyms=["orders model"], instructions="Use for revenue and customer analytics.", ), ) return model, [orders_ds, customers_ds, total_revenue, double_revenue] def main() -> None: model, entities = build_graph() all_mcps: list[MetadataChangeProposalWrapper] = [] all_mcps.extend(model.as_mcps()) for entity in entities: all_mcps.extend(entity.as_mcps()) records: list[dict[str, Any]] = [dict(mcp.to_obj()) for mcp in all_mcps] with open("semantic_model_create.json", "w") as f: json.dump(records, f, indent=2, default=str) print(f"Wrote {len(all_mcps)} MCPs to semantic_model_create.json") # 向真实服务端发出时,建议先调用 opt-in 预检助手(见下文): # from datahub.sdk import DataHubClient, require_metrics_support # client = DataHubClient(server=..., token=...) # require_metrics_support(client) # raises if the server version is too old # for entity in [model, *entities]: # client.entities.upsert(entity) if __name__ == "__main__": main()

要点说明:

  • 逻辑数据集的name建议编码为<sm_path>.<sm_id>.<view_name>(例如analytics.orders_model.orders_ds),保证逻辑数据集在不同 semantic model 之间保持唯一——这一点由 SemanticModelDataset 源码 的文档明确约定;
  • 每个SemanticFieldInput通过field_path+type定义字段,semantic_typeDIMENSIONMEASURE,可选地携带expression(多语言方言表达式)、aggregation_functionis_time_dimension以及用于 AI 场景的ai_context
  • Metricexpression支持三种输入形态:纯字符串(默认按 ANSI SQL 方言)、单个DialectExpressionInput、或DialectExpressionInput列表;double_revenue展示了不写expression、仅通过derived_from派生自另一个指标的写法。

SDK 替你发出了什么

对每个 builder 调用entity.as_mcps()时,SDK 会生成完整的 aspect 集合并自动装配血缘链(实现见 semantic_model.py 与 metric.py):

  • semanticModel:一个Status、一个SemanticModelInfo(name、description、可选relationships——成员关系不在这里列出),以及模型级AiContext仅在非空时发出)。从源码看,成员关系只存在于成员一侧:每个逻辑数据集通过semanticModelProperties.semanticModelIsPartOf)指向模型,semanticModelInfo上没有datasets/metrics数组([SemanticModel 类文档](https://link.gitcode.com/i/9705b1b974a942664653552b300a072e#L127-L157))。
  • 逻辑datasets:每个都会带上SubTypes([SEMANTIC_MODEL_DATASET])、一个SemanticModelProperties(alias, semanticModel=<model urn>)成员指针(IsPartOf)、一个包含已声明字段的SchemaMetadata,以及——当提供了upstreams时——指向物理数据集的UpstreamLineage。对每个字段,SDK 发出锚定在schemaFieldURN 上的semanticFieldAnnotation(未提供expression时自动合成为f"{alias}.{field_path}"),并在非空时发出字段锚定的aiContext
  • metrics:每个都会带上StatusMetricInfo(含semanticModel=<model urn>成员关系(ModeledBy)与可选expression;省略时绝不伪造表达式)、MetricRelationships总是发出,即使derivedFrom为空,这样hasParentMetric才能索引为 false)、可选的MetricUpstreams.datasetUpstreams(指向指标读取的 Semantic Model Dataset URN),以及仅在非空时的AiContext

血缘经由metricUpstreams(Metric → 逻辑数据集)和每个逻辑数据集的upstreamLineage(逻辑数据集 → 物理数据集)形成Metric → Logical Dataset → Physical Dataset。SemanticModel 是成员的容器(bounding box),不是血缘跳。

值得注意的源码细节:

  • SemanticModel.as_mcps()在发出时会以strict=True重新校验 relationships:join 的 alias 必须匹配某个已挂载逻辑数据集的 alias,join 列必须出现在该数据集 schema 中,列数两侧必须相等;别名重复或结构非法时抛出SdkUsageError(校验逻辑)。
  • SemanticModelDataset._new_from_graph在从服务端读取时绕过严格的__init__,字段注解被标记为 create-only,读取构造的实例不携带任何字段注解(读取路径)。

期望输出

运行示例后,当前工作目录会写出semantic_model_create.json。打开它并核对 aspect 形状是否符合生产者契约:

  • URN 模式:urn:li:semanticModel:(urn:li:dataPlatform:snowflake,analytics,orders_model)urn:li:metric:(urn:li:dataPlatform:snowflake,analytics,total_revenue)urn:li:dataset:(urn:li:dataPlatform:snowflake,analytics.orders_model.orders_ds,PROD)
  • 逻辑数据集带有Semantic Model Dataset子类型;
  • 每个逻辑数据集的semanticModelProperties以正确的alias指回模型 URN;
  • semanticFieldAnnotationMCP 锚定在schemaFieldURN 上,expression在未显式提供时回退为ORDERS.order_id
  • aiContext只出现在输入非空的字段/实体上(build_ai_context在四个字段全空时返回None,见 _semantic_shared.py);
  • total_revenue发出metricUpstreams.datasetUpstreams,指向orders_ds的 Semantic Model Dataset URN(来自upstream_datasets=[orders_ds.urn])。

API 参考

每个 builder 的完整参数面见 SDK Entities Reference(python-sdk/sdk-v2/entities.mdx,位于仓库 python-sdk 文档目录)。核心类:

  • SemanticModeldatahub.sdk.semantic_model.SemanticModel
  • SemanticModelDatasetdatahub.sdk.semantic_model.SemanticModelDataset
  • Metricdatahub.sdk.metric.Metric

常用输入类型(均可从datahub.sdk顶层导入):

输入类型用途说明
SemanticFieldInput逻辑数据集的字段field_pathtypesemantic_type必填;expression省略时自动合成f"{alias}.{field_path}"
SemanticModelRelationshipInput数据集间 join 路径from_alias/to_alias必须与逻辑数据集的alias匹配
DialectExpressionInput(dialect, expression) 对可单用、列表使用或包裹在MetricExpressionInputType
AiContextInputaiContextaspect 输入synonyms/instructions/examples/custom_instructions,全空则不发 aspect

服务端兼容性

semanticModelmetric与逻辑dataset实体要求服务端 build 已注册 semantic-model 元数据模型。向未注册这些 aspect 的服务端发出会响亮地失败——服务端拒绝未注册的 aspect,emit_mcps抛出异常。

如果你想在服务端拒绝之前拿到一个清晰、可操作的错误,可在发出前调用 opt-in 的预检助手:

from datahub.sdk import DataHubClient, require_metrics_support client = DataHubClient(server="...", token="...") require_metrics_support(client) # raises if the server version is too old

该助手委托给RestServiceConfig.supports_feature:当服务端报告的版本不支持这些实体时抛出异常;当没有可检查的版本信号时fail open(此时由运维负责确保运行包含该模型的 build)。它没有被自动织入DataHubClient.upsert——想要预检时需显式调用。

从 实现源码 看,其行为更细致一些:

  • 接受DataHubClient或原始DataHubGraph,会自动解包 client 内部的 graph;
  • 特性检查使用ServiceFeature.SEMANTIC_MODEL_ENTITIES;非 semver build(dev/snapshot/sha tag)会抛ValueError,此时 fail open 而不泄漏原始解析错误;
  • 仅当服务端报告版本低于最低要求、且是托管(Cloud)服务端时才抛SdkUsageError——OSS/自托管部署不被 SDK 做版本门槛限制,避免阻塞所有 OSS 发出路径。

逻辑数据集的读-改-写注意事项

逻辑数据集SemanticModelDataset上的逐字段semanticFieldAnnotation和字段级aiContextcreate-only的:逻辑数据集与dataset共用实体类型,因此client.entities.get(<dataset urn>)会把它作为基础Dataset水合——锚定在字段上的注解存在于schemaFieldURN 上,而非 dataset 的 aspect 袋中,读取时不会被带回。要更新一个逻辑数据集,应重新构建一个新的SemanticModelDataset,通过schema构造参数重新挂载字段,而不是对读回的Dataset做读-改-写。

【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/16 15:19:58

自指意义文明:认知科学与技术架构解析

1. 概念解析&#xff1a;什么是"自指意义文明"&#xff1f;"自指意义文明"这个复合概念由三个关键词构成&#xff1a;自指、意义、文明。拆解来看&#xff0c;"自指"指的是系统能够反身指向自身&#xff0c;形成递归式的认知结构&#xff1b;&qu…

作者头像 李华
网站建设 2026/9/16 15:18:30

自研轻量级全景监控系统:地图渲染与实时推送实战

运维和数据产品这个圈子里&#xff0c;有一种需求几乎每个团队都会遇到&#xff1a;设备分散在园区各个角落&#xff0c;业务系统各自为战&#xff0c;想看一眼全局状态&#xff0c;得同时打开七八个后台&#xff1b;真出了故障&#xff0c;排查链路基本靠电话和口头确认&#…

作者头像 李华
网站建设 2026/9/16 15:17:35

Niushop v5.1.7电商源码:LNMP部署、多模版切换与支付回调实战解析

简介&#xff1a;这是一套基于 ThinkPHP6 的 Niushop 多模板大型商城电商系统源码 v5.1.7&#xff0c;主要面向需要快速搭建或二次开发网上商城的企业、团队与 PHP 开发者。系统覆盖普通商品与虚拟商品管理、二维码核销、拼团/分销/积分兑换等营销玩法&#xff0c;支持物流配送…

作者头像 李华
网站建设 2026/9/16 15:14:28

AT89C52停车场车位管理系统设计:红外检测与PWM舵机控制

简介&#xff1a;面向单片机课程设计与毕业设计的AT89C52停车场车位管理系统方案包&#xff0c;模拟16车位管理场景&#xff0c;通过按键模拟车辆进出、LED状态指示与LCD实时显示&#xff0c;满位时触发报警&#xff0c;可帮助学习者快速掌握外围器件驱动与中断/按键逻辑设计。…

作者头像 李华
网站建设 2026/9/16 15:12:55

CloudStream实操教程:5分钟搭好你的专属流媒体中心

CloudStream实操教程&#xff1a;5分钟搭好你的专属流媒体中心 【免费下载链接】cloudstream Android app for streaming and downloading media. 项目地址: https://gitcode.com/GitHub_Trending/cl/cloudstream 手机里装了五个播放器&#xff0c;每个都带广告&#xf…

作者头像 李华