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_type取DIMENSION或MEASURE,可选地携带expression(多语言方言表达式)、aggregation_function、is_time_dimension以及用于 AI 场景的ai_context; Metric的expression支持三种输入形态:纯字符串(默认按 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.semanticModel(IsPartOf)指向模型,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:每个都会带上Status、MetricInfo(含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 文档目录)。核心类:
SemanticModel—datahub.sdk.semantic_model.SemanticModelSemanticModelDataset—datahub.sdk.semantic_model.SemanticModelDatasetMetric—datahub.sdk.metric.Metric
常用输入类型(均可从datahub.sdk顶层导入):
| 输入类型 | 用途 | 说明 |
|---|---|---|
SemanticFieldInput | 逻辑数据集的字段 | field_path、type、semantic_type必填;expression省略时自动合成f"{alias}.{field_path}" |
SemanticModelRelationshipInput | 数据集间 join 路径 | from_alias/to_alias必须与逻辑数据集的alias匹配 |
DialectExpressionInput | (dialect, expression) 对 | 可单用、列表使用或包裹在MetricExpressionInputType中 |
AiContextInput | aiContextaspect 输入 | synonyms/instructions/examples/custom_instructions,全空则不发 aspect |
服务端兼容性
semanticModel、metric与逻辑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和字段级aiContext是create-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),仅供参考