DataHub APIs 与 SDK 全览:在 GraphQL、OpenAPI、Python/Java SDK 与 CLI 之间做出正确选择
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
本篇技术指南以 docs/api/datahub-apis.md 为核心,系统梳理 DataHub 面向元数据管理的全部接口家族——GraphQL API、OpenAPI v3 接口、Python/Java SDK 以及底层事件写入通道,并给出官方推荐的选型原则、完整的能力对比矩阵、异步写入的验证边界与可运行的实操示例。读完本文,你将能够根据自身用例(UI 驱动的交互查询、批量元数据注入、自定义实体建模、程序化管线)快速锁定正确的接口,并掌握各接口的认证、容错与条件写入等关键细节。
一、DataHub API 全景:四种接口的定位与取舍
DataHub 在平台之上提供了多套用于操作元数据的接口。核心关联文档给出了一张直观的选型表,我们在此基础上补充了实现细节与仓库佐证,汇总如下:
| API | 定义 | 优势 | 局限 |
|---|---|---|---|
| Python SDK | SDK | 高度灵活,适合批量执行 | 需要理解元数据变更事件(MCP) |
| Java SDK | SDK | 高度灵活,适合批量执行 | 需要理解元数据变更事件(MCP) |
| GraphQL API | GraphQL 接口 | 直观,与 UI 能力镜像对齐 | 灵活性低于 SDK;需要掌握 GraphQL 语法 |
| OpenAPI | 面向高级用户的底层 API | 最强大、最灵活 | 对简单用例而言上手门槛偏高;无配套 SDK,但产品内会生成 OpenAPI 规范 |
总体而言,Python 与 Java SDK 是官方最推荐用于扩展和定制 DataHub 实例行为的工具,尤其适合程序化使用场景。这一点在文档中多次被强调,也是后续各节选型讨论的基调。
值得补充的是,除了上表中的四类接口,DataHub 还保留了经典的Rest.li API(如ingestProposal),它是 GMS 最早期的写入通道,SDK 底层仍会与其兼容交互。此外,官方文档还提示用户关注 CLI 命令(如 docs/cli-commands/dataset.md)作为运维侧的补充手段。
SDK:官方推荐的首选程序化方案
Python 与 Java SDK 提供了完整的 CRUD 能力以及构建在 DataHub 之上的任意复杂功能,官方建议大多数用例直接使用 SDK。常见的高价值场景包括:
- 在数据实体之间定义血缘(lineage)关系;
- 执行批量操作,例如为多个数据集批量打标签;
- 创建自定义元数据实体。
SDK 的本质是围绕MetadataChangeProposal(MCP)事件构建的发射器(Emitter)封装。以 Python 为例,acryl-datahub包同时提供REST Emitter与Kafka Emitter两套 API,分别面向“确认写入”与“高吞吐解耦”两种诉求:
- REST Emitter:基于
requests模块的轻量封装,提供阻塞式接口,适合需要确认元数据已持久化、以及存在“写后读”场景的使用方式; - Kafka Emitter:基于
confluent-kafka的SerializingProducer,提供非阻塞接口,适合希望利用 Kafka 作为高可用消息总线、在 DataHub 元数据服务宕机时仍能持续采集元数据的场景。需要特别注意的是,Kafka Emitter 使用Avro序列化元数据事件,更换序列化器将导致事件不可被处理。
Python 端两个 Emitter 的核心实现分别位于 rest_emitter.py 与 kafka_emitter.py。Java 端则由io.acryl:datahub-client包提供 REST、Kafka 与 File 三类 Emitter,其中 File Emitter 可将 MCP 写入 JSON 文件,随后通过 Metadata File source 离线导入,适用于生产系统无法直连 GMS 或 Kafka 的场景。
Python REST Emitter 实操示例
import datahub.emitter.mce_builder as builder from datahub.emitter.mcp import MetadataChangeProposalWrapper from datahub.metadata.schema_classes import DatasetPropertiesClass from datahub.emitter.rest_emitter import DatahubRestEmitter # 创建指向 DataHub GMS 的 REST Emitter emitter = DatahubRestEmitter(gms_server="http://localhost:8080", extra_headers={}) # DataHub Cloud 场景则指向托管域的 GMS 端点并携带 token # emitter = DatahubRestEmitter(gms_server="https://<your-domain>.acryl.io/gms", token="<your token>", extra_headers={}) # 测试连接 emitter.test_connection() # 构造数据集属性对象 dataset_properties = DatasetPropertiesClass( description="This table stored the canonical User profile", customProperties={"governance": "ENABLED"}, ) # 构造 MetadataChangeProposalWrapper 对象 metadata_event = MetadataChangeProposalWrapper( entityUrn=builder.make_dataset_urn("bigquery", "my-project.my-dataset.user-table"), aspect=dataset_properties, ) # 阻塞式发射元数据 emitter.emit(metadata_event)发射模式(Emit Mode):吞吐与一致性的权衡
emit()与emit_mcp()接受可选的emit_mode: EmitMode参数,控制调用等待时间与一致性保证。Emitter 默认使用SYNC_PRIMARY,可逐次调用覆盖,也可在 Emitter 上一次性设置default_emit_mode。
EmitMode | 调用返回时机 | 适用场景 |
|---|---|---|
SYNC_WAIT | SQL与Elasticsearch 均已更新 | 写入后必须立即可搜索,一致性最强但最慢 |
SYNC_PRIMARY(默认) | SQL 已更新(ES 异步索引) | 低量写入且需要“写后读”直接实体获取 |
ASYNC | 变更已入队(立即返回) | 高吞吐或批量摄取,可容忍最终一致 |
ASYNC_WAIT | 已确认排队变更持久化 | 希望异步批处理/并行,但仍需持久化确认 |
选型建议:低量写入且需要回读或调用点即时报错 → 保持SYNC_PRIMARY(若写入后必须立即可搜索则用SYNC_WAIT);高吞吐或批量摄取 → 显式设置为ASYNC。默认的SYNC_PRIMARY并不适合高吞吐场景——每次写入都做同步主存储提交会给 GMS 及其底层 SQL 存储带来沉重负载。
使用ASYNC需注意两个后果:其一,emit()不会对被拒绝或无效的写入抛出异常——调用在变更入队后即成功返回,校验或持久化失败会稍后出现在 Failed-MCP topic 与消费端日志中;其二,不保证“写后读”,任何“先写后立即读同一实体”的流程都必须容忍最终一致性。
Java SDK 快速上手
Java SDK(V1)以io.acryl:datahub-client提供,Gradle 依赖声明如下:
implementation 'io.acryl:datahub-client:__version__'Maven 方式则添加:
<dependency> <groupId>io.acryl</groupId> <artifactId>datahub-client</artifactId> <version>__version__</version> </dependency>REST Emitter 采用 lambda 风格的可变构造器模式,配置参数与 Python Emitter 大体镜像:
RestEmitter emitter = RestEmitter.create(b -> b .server("http://localhost:8080") // Auth token for DataHub Cloud // .token(AUTH_TOKEN_IF_NEEDED) // Override default timeout of 10 seconds // .timeoutSec(OVERRIDE_DEFAULT_TIMEOUT_IN_SECONDS) // Add additional headers // .extraHeaders(Collections.singletonMap("Session-token", "MY_SESSION")) );发射 MCP 支持阻塞式(Future.get())与回调式两种写法,官方文档的完整示例见 as-a-library.md。值得注意的是,仓库中 Java SDK 相关文档已推荐新项目优先使用Java SDK V2(as-a-library-v2.md),它提供类型安全的实体构造器、简化的 CRUD 操作与基于 Patch 的高效元数据更新。
二、GraphQL API:与 UI 能力对齐的高级查询接口
graphqlAPI 是 DataHub 前端所使用的主 API。它默认带有缓存、同步操作以及其他面向 UI 的预期行为,因此在程序化抓取与更新时需要谨慎——其操作在范围上被有意限制,定位是简化最常见操作的高级 API。
GraphQL API 很适合刚接触 DataHub 的用户,尤其是配合 GraphiQL 使用时更友好、更直接。典型用例包括:
- 带条件地搜索数据集;
- 查询实体之间的关系。
查询(Queries):读取实体
以下 GraphQL 查询获取指定数据集的urn与properties.name:
{ dataset(urn: "urn:li:dataset:(urn:li:dataPlatform:kafka,SampleKafkaDataset,PROD)") { urn properties { name } } }除 URN 与属性外,还可获取某资产的所有者、标签、域、术语等元数据。相关查询指南包括:
- 查询数据集的所有者
- 查询数据集的标签
- 查询数据集的域
- 查询数据集的术语
- 查询数据集的弃用状态
- 查询 DataFlow 下所有 DataJob
搜索(Search):全文检索
使用search(input: SearchInput!)查询对特定类型实体执行全文检索:
{ search(input: { type: DATASET, query: "my sql dataset", start: 0, count: 10 }) { start count total searchResults { entity { urn type ...on Dataset { name } } } } }input参数指定实体类型、查询词、起始索引与返回数量。query支持通配模式:
*:搜索全部实体;*[string]:搜索方面(aspect)以指定字符串开头的实体;[string]*:搜索方面以指定字符串结尾的实体;*[string]*:搜索方面匹配指定字符串的实体;[string]:搜索方面包含指定字符串的实体。
:::note 默认情况下 Elasticsearch 只允许通过 search API 分页浏览 10,000 个实体。如需更多,可调整 ES 的index.max_result_window配置,或使用 scroll API 直接读取索引。 :::
变更(Mutations):更新实体
变更实体元数据的 Mutations 受 DataHub Access Policies 约束,服务端会校验请求方 actor 是否被授权执行该操作。同时官方明确提示:GraphQL mutations 主要面向 UI 交互设计,程序化用例应避免使用——它们不适合数据集成工作流中的高吞吐或批量场景。程序化元数据管理、数据摄取与批量操作请使用 Python SDK。
以更新 Dashboard 实体为例:
mutation updateDashboard { updateDashboard( urn: "urn:li:dashboard:(looker,baz)", input: { editableProperties: { description: "My new description" } } ) { urn } }更多变更操作指南:
- 添加标签 / 移除标签
- 添加术语 / 移除术语
- 添加域 / 移除域
- 添加所有者 / 移除所有者
- 更新弃用状态
- 编辑数据集/列的描述文档
- 软删除
错误处理:检查 data 与 errors 两个字段
GraphQL 请求出错时并不总是返回非 200 的 HTTP 响应体。错误会出现在响应体顶层的errors字段中,这允许客户端优雅地处理应用服务器返回的部分数据。因此每次请求后必须同时检查data与errors字段。错误对象包含 message、path 以及携带标准错误码的 extensions:
{ "errors": [ { "message": "Failed to change ownership for resource urn:li:dataFlow:(airflow,dag_abc,PROD). Expected a corp user urn.", "locations": [ { "line": 1, "column": 22 } ], "path": ["addOwners"], "extensions": { "code": 400, "type": "BAD_REQUEST", "classification": "DataFetchingException" } } ] }官方支持的错误码如下:
| Code | Type | Description |
|---|---|---|
| 400 | BAD_REQUEST | 查询或变更格式错误 |
| 403 | UNAUTHORIZED | 当前 actor 未被授权执行请求的操作 |
| 404 | NOT_FOUND | 资源不存在 |
| 500 | SERVER_ERROR | 发生内部错误,请检查服务器日志或联系 DataHub 管理员 |
三、OpenAPI:面向高级用户的底层 REST 接口
OpenAPI 标准是广泛使用的 REST 风格 API 文档化与设计方法。DataHub 基于此发布了一套 OpenAPI 端点,方便第三方系统深度集成。详细的使用指南见 openapi-usage-guide.md。
定位 OpenAPI 端点
OpenAPI 端点当前隔离在 GMS 上的一个独立 Servlet 中,随 GMS 自动部署。该 Servlet 内建自动生成的 OpenAPI UI(即 Swagger),访问路径为GMS_SERVER_HOST:GMS_PORT/openapi/swagger-ui/index.html,本地 Quickstart 对应 http://localhost:8080/openapi/swagger-ui/index.html。
前端也会以代理方式暴露同一端点(将 GMS 主机与端口替换为前端 URL),本地 Quickstart 对应 http://localhost:9002/openapi/swagger-ui/index.html,并可在用户头像右上角下拉菜单中直接打开。
原始 JSON/YAML 格式的 OpenAPI 规范可通过BASE_URL/openapi/v3/api-docs或BASE_URL/openapi/v3/api-docs.yaml获取,可喂给 codegen 系统生成任意语言的客户端代码(不同语言 codegen 成熟度不一,可能需要定制)。UI 中的请求/响应对象 Schema 均在构建期由 PDL 模型自动生成。
从仓库源码看,metadata-service/openapi-servlet模块中的 SpringWebConfig.java 按包名划分了 v1/v2 等不同版本的端点集合,印证了 OpenAPI Servlet 是随 GMS 统一装配的独立 REST 层。
主要端点分类
| 端点 | 用途 |
|---|---|
/entities | 对元数据图进行读写。整个 DataHub 元数据模型都可以以“实体 + 方面(aspect)”的形式写入,或按需读取单个实体的元数据 |
/relationships | 查询图结构,从一个实体导航到其他实体的关系 |
/timeline | 查询指定实体的版本化历史,例如数据集的所有 Schema 变更或文档变更记录,详见 timeline 指南 |
/platform | 更底层的 API,允许以标准格式将元数据事件写入 DataHub 平台 |
实体端点实操:UPSERT / CREATE / GET / DELETE
UPSERT(POST):不带额外 URL 参数的 POST 执行方面 UPSERT,实体不存在则创建、存在则更新:
curl --location --request POST 'localhost:8080/openapi/entities/v1/' \ --header 'Content-Type: application/json' \ --header 'Accept: application/json' \ --header 'Authorization: Bearer <token>' \ --data-raw '[ { "aspect": { "__type": "SchemaMetadata", "schemaName": "SampleHdfsSchema", "platform": "urn:li:dataPlatform:platform", "platformSchema": { "__type": "MySqlDDL", "tableSchema": "schema" }, "version": 0, "created": { "time": 1621882982738, "actor": "urn:li:corpuser:etl", "impersonator": "urn:li:corpuser:jdoe" }, "lastModified": { "time": 1621882982738, "actor": "urn:li:corpuser:etl", "impersonator": "urn:li:corpuser:jdoe" }, "hash": "", "fields": [ { "fieldPath": "county_fips_codefg", "jsonPath": "null", "nullable": true, "description": "null", "type": { "type": { "__type": "StringType" } }, "nativeDataType": "String()", "recursive": false }, { "fieldPath": "county_name", "jsonPath": "null", "nullable": true, "description": "null", "type": { "type": { "__type": "StringType" } }, "nativeDataType": "String()", "recursive": false } ] }, "entityType": "dataset", "entityUrn": "urn:li:dataset:(urn:li:dataPlatform:platform,testSchemaIngest,PROD)" } ]'CREATE(POST +createEntityIfNotExists=true):仅当实体不存在时写入,实体已存在则返回错误而非覆盖:
curl --location --request POST 'localhost:8080/openapi/entities/v1/?createEntityIfNotExists=true' \ --header 'Content-Type: application/json' \ --header 'Accept: application/json' \ --header 'Authorization: Bearer <token>' \ --data-raw '<see previous example>'实体已存在时返回如下 422 错误:
422 ValidationExceptionCollection{EntityAspect:(urn:li:dataset:(urn:li:dataPlatform:platform,testSchemaIngest,PROD),schemaMetadata) Exceptions: [com.linkedin.metadata.aspect.plugins.validation.AspectValidationException: Cannot perform CREATE if not exists since the entity key already exists.]}
GET(读取最新方面):
curl --location --request GET 'localhost:8080/openapi/entities/v1/latest?urns=urn:li:dataset:(urn:li:dataPlatform:platform,testSchemaIngest,PROD)&aspectNames=schemaMetadata' \ --header 'Accept: application/json' \ --header 'Authorization: Bearer <token>'DELETE(软删除):
curl --location --request DELETE 'localhost:8080/openapi/entities/v1/?urns=urn:li:dataset:(urn:li:dataPlatform:platform,testSchemaIngest,PROD)&soft=true' \ --header 'Accept: application/json' \ --header 'Authorization: Bearer <token>'其中soft=true(默认)表示软删除,soft=false则为硬删除。官方使用指南中还附带了一份完整的 Postman Collection(含 POST/GET/DELETE 与 SchemaMetadata 方面的示例),可直接导入调试。
关系端点实操
GET示例——查询与用户datahub具有IsPartOf入向关系的实体:
curl -X 'GET' \ 'http://localhost:8080/openapi/relationships/v1/?urn=urn%3Ali%3Acorpuser%3Adatahub&relationshipTypes=IsPartOf&direction=INCOMING&start=0&count=200' \ -H 'accept: application/json'示例响应:
{ "start": 0, "count": 2, "total": 2, "entities": [ { "relationshipType": "IsPartOf", "urn": "urn:li:corpGroup:bfoo" }, { "relationshipType": "IsPartOf", "urn": "urn:li:corpGroup:jdoe" } ] }程序化使用:Java Rest Emitter
OpenAPI 模型的程序化使用可通过包含生成模型的 Java Rest Emitter 完成。最小 Java 项目需要以下依赖(Gradle 格式):
dependencies { implementation 'io.acryl:datahub-client:<DATAHUB_CLIENT_VERSION>' implementation 'org.apache.httpcomponents:httpclient:<APACHE_HTTP_CLIENT_VERSION>' implementation 'org.apache.httpcomponents:httpasyncclient:<APACHE_ASYNC_CLIENT_VERSION>' }通过构造UpsertAspectRequest列表向/platform/entities/v1端点批量发射元数据事件:
import io.datahubproject.openapi.generated.DatasetProperties; import datahub.client.rest.RestEmitter; import datahub.event.UpsertAspectRequest; import java.io.IOException; import java.util.ArrayList; import java.util.List; import java.util.concurrent.ExecutionException; public class Main { public static void main(String[] args) throws IOException, ExecutionException, InterruptedException { RestEmitter emitter = RestEmitter.createWithDefaults(); List<UpsertAspectRequest> requests = new ArrayList<>(); UpsertAspectRequest upsertAspectRequest = UpsertAspectRequest.builder() .entityType("dataset") .entityUrn("urn:li:dataset:(urn:li:dataPlatform:bigquery,my-project.my-other-dataset.user-table,PROD)") .aspect(new DatasetProperties().description("This is the canonical User profile dataset")) .build(); UpsertAspectRequest upsertAspectRequest2 = UpsertAspectRequest.builder() .entityType("dataset") .entityUrn("urn:li:dataset:(urn:li:dataPlatform:bigquery,my-project.another-dataset.user-table,PROD)") .aspect(new DatasetProperties().description("This is the canonical User profile dataset 2")) .build(); requests.add(upsertAspectRequest); requests.add(upsertAspectRequest2); System.out.println(emitter.emit(requests, null).get()); System.exit(0); } }四、OpenAPI v3 进阶特性:条件写入、批量读取与通用 Patch
条件写入(Conditional Writes)
所有方面的 create/POST 端点都在 POST body 中支持headers,以支撑批量 API。这些 header 用于实现条件写入语义,详细文档见 MetadataChangeProposal 与 MetadataChangeLog 事件。其核心机制包括:
If-Version-Match:方面每次更新都会递增一个version并存入SystemMetadata。写入方可在请求中携带期望版本,若不匹配则写入失败,从而防止覆盖其他进程已修改的方面。若方面尚不存在,其版本为-1,可借此实现“仅创建”语义。If-Modified-Since/If-Unmodified-Since:基于时间的条件写入,日期须符合 ISO-8601 标准,防止目标方面在读取后被修改时仍执行写入。- Change Type
CREATE/CREATE_ENTITY:分别表示“方面不存在才创建”与“实体无任何方面才创建”。默认违反约束会抛出校验异常;若希望丢弃该写入而不视为异常,可附加 headerIf-None-Match: *。
批量读取(Batch Get)
所有实体都存在/v3/entity/{entityName}/batchGet形式的批量读取端点,可一次获取实体及其多个方面,默认返回最新版本;结合If-Version-Matchheader 可检索指定版本。该接口目前每个实体/方面仅返回单一版本,但不同实体之间可指定不同版本。
示例请求——携带systemMetadata=true以查看方面当前版本:
[ { "urn": "urn:li:dataset:(urn:li:dataPlatform:hive,fct_users_deleted,PROD)", "globalTags": {}, "datasetProperties": {} }, { "urn": "urn:li:dataset:(urn:li:dataPlatform:hive,fct_users_created,PROD)", "globalTags": {}, "datasetProperties": {} } ]响应中systemMetadata会携带每个方面的"version"。当对第二个 URN 的globalTags追加一个新标签后,该方面的版本会从1递增为2;随后使用If-Version-Match: 1的 headers 即可回读前一版本。完整往返示例见 openapi-usage-guide.md。
通用 Patching:基于 RFC 6902 的数组主键扩展
OpenAPI v3 的 PATCH 端点消除了以往 Patch 对后端专用代码的依赖,它基于 JSON Patch 标准(RFC 6902),并针对数组操作做了显著增强。为保持向后兼容,默认仍使用传统 Patch 模板;当arrayPrimaryKeys非空或forceGenericPatch为true时启用通用 Patch。
标准 JSON Patch 对数组只支持基于索引的操作(add/[index]、remove/[index]、replace/[index]),在数组顺序不可预测、多客户端并发修改或客户端不知道当前索引时会产生问题。DataHub 通过arrayPrimaryKeys字段将数组操作转换为类 Map 操作:
- 数组在概念上被视为 Map,每个元素可由其主键寻址;
- 主键可以是复合的(多个字段组合);
- 路径表达式使用这些键而非数字索引;
- 后端负责在 Map 式操作与实际数组修改之间转换。
这带来了幂等(与数组顺序无关)、定向修改(无需知道当前数组状态)、并发无冲突(修改不同数组元素时)以及更直观的 API 使用体验。
主键定义与路径构造示例——为globalTags方面添加带来源(attribution)归属的标签:
{ "arrayPrimaryKeys": { "tags": ["attribution␟source", "tag"] }, "patch": [ { "op": "add", "path": "/tags/urn:li:platformResource:source1/urn:li:tag:tag1", "value": { "tag": "urn:li:tag:tag1", "attribution": { "source": "urn:li:platformResource:source1", "actor": "urn:li:corpuser:user", "time": 0 } } } ] }其中tags是被 Patch 的数组字段;主键是attribution.source与tag的复合;␟(Unit Separator,U+241F)分隔符表示第一个键分量中的嵌套路径。路径/tags/urn:li:platformResource:source1/urn:li:tag:tag1中,系统将urn:li:platformResource:source1作为attribution.source的值、urn:li:tag:tag1作为 tag 的值,据此匹配数组元素。
当前实现支持的标准 JSON Patch 操作:
| Operation | Description |
|---|---|
| add | 添加新元素;若键匹配的元素已存在则替换 |
| remove | 按键移除匹配的元素 |
remove操作示例:{"op": "remove", "path": "/tags/urn:li:platformResource:source1/urn:li:tag:tag1"},它会仅移除匹配复合键的元素,即使其他元素存在部分匹配的键也会被保留。
五、深入理解 MCP/MCL:所有写入通道的共同底层
无论是 SDK、OpenAPI 还是 Rest.li,所有元数据写入最终都会落到MetadataChangeProposal(MCP)事件流上,由 GMS 校验后产生MetadataChangeLog(MCL)供下游消费。理解这一层是“为什么 SDK 需要理解元数据变更事件”这一局限的根源,其完整模型见 mcp-mcl.md。
MCP 的核心字段包括:entityType(实体类型,如 dataset、chart)、entityUrn(实体 URN,与entityKeyAspect二选一)、changeType(变更类型:UPSERT、CREATE、CREATE_ENTITY、UPDATE、DELETE、PATCH,当前仅支持前五者与PATCH)、aspectName、aspect(GenericAspect,含value与contentType,目前仅支持application/json)以及systemMetadata与headers。写入的原子单元是单个方面。
异步 GMS 写入的验证边界(官方重点警告)
关联文档特别强调了一个容易混淆的关键点——GMS 的异步写入(async=true)与直连 Kafka 生产事件有着本质区别:
- GMS 异步接受 ≠ 跳过校验:当你通过 GMS API(Rest.li
ingestProposal、OpenAPI 实体写入,以及针对这些端点的 SDK 客户端)以async=true提交 MCP 时,GMS 会在接受请求之前运行完整的提案校验管线,包括 Schema 检查、实体级授权(isAPIAuthorized)以及已注册的方面负载校验器(例如标签权限约束、logicalParent等特定方面的授权)。未授权或无效的提案会以 403/422 被同步拒绝,不会发布到 Kafka。 - 异步仅指主存储提交延迟:
async=true只表示 GMS 不在接受时把写入提交到主存储,而是将 MCP 发布到MetadataChangeProposaltopic;之后由 MCE consumer 以async=false应用。此阶段发生的失败(如预提交校验或存储错误)会进入Failed MCP topic,不会回传给最初的 API 调用方。 - 不要混淆异步 GMS 摄取与直接 Kafka 生产:直接向
MetadataChangeProposaltopic 写入 MCP 会绕过GMS 接受时的授权与校验。MCE consumer 在系统上下文下处理这些消息,并非第二道用户授权闸门。因此必须对 Kafka 的访问权限加以限制。
这一边界也解释了 Python SDK 中ASYNC发射模式的语义:调用成功只代表“变更已入队”,校验/持久化失败将稍后在 Failed-MCP topic 与消费端日志中显现。
条件写入与同步索引更新
除上述If-Version-Match、If-Modified-Since等条件写入 header 外,MCP 还支持:
X-DataHub-Sync-Index-Update: true:默认 Elasticsearch 更新是异步的,添加该 header 可对特定 MCP 启用同步索引更新;- 方面大小校验:可通过环境变量与
datahub.validation.aspectSize配置(prePatch/postPatch两阶段,warnSizeBytes仅告警不阻断,maxSizeBytes默认 16MB 并配合IGNORE/DELETE补救策略),用于防护超大方面的写入。
相关 Kafka Topic 一览
| Topic | 作用 |
|---|---|
MetadataChangeProposal_v1 | 异步摄取到 GMS 的提案流;失败会产出 Failed MCP |
FailedMetadataChangeProposal_v1 | 摄取失败的提案 |
MetadataChangeLog_Versioned_v1 | 版本化方面的变更日志(默认保留 7 天) |
MetadataChangeLog_Timeseries_v1 | 时间序列方面的变更日志(默认保留 90 天,可回放备份) |
六、DataHub API 能力对比矩阵
下表源自关联文档的官方对比(最后更新:2024-02-16),覆盖 GraphQL、Python SDK 与 OpenAPI 三类接口在常见元数据操作上的能力差异。最显著的模式是:凡是涉及“创建/新增”的操作,GraphQL 大多不支持(🚫),而 Python SDK 与 OpenAPI 均支持(✅)——这与“GraphQL 面向 UI 交互、SDK/OpenAPI 面向程序化写入”的定位完全一致。
| Feature | GraphQL | Python SDK | OpenAPI |
|---|---|---|---|
| Create a Dataset | 🚫 | ✅ [Guide] | ✅ |
| Delete a Dataset (Soft Delete) | ✅ [Guide] | ✅ [Guide] | ✅ |
| Delete a Dataset (Hard Delete) | 🚫 | ✅ [Guide] | ✅ |
| Search a Dataset | ✅ [Guide] | ✅ | ✅ |
| Read a Dataset Deprecation | ✅ | ✅ | ✅ |
| Read Dataset Entities (V2) | ✅ | ✅ | ✅ |
| Create a Tag | ✅ [Guide] | ✅ [Guide] | ✅ |
| Read a Tag | ✅ [Guide] | ✅ [Guide] | ✅ |
| Add Tags to a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Add Tags to a Column of a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Remove Tags from a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Create Glossary Terms | ✅ [Guide] | ✅ [Guide] | ✅ |
| Read Terms from a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Add Terms to a Column of a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Add Terms to a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Create Domains | ✅ [Guide] | ✅ [Guide] | ✅ |
| Read Domains | ✅ [Guide] | ✅ [Guide] | ✅ |
| Add Domains to a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Remove Domains from a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Create / Upsert Users | ✅ [Guide] | ✅ [Guide] | ✅ |
| Create / Upsert Group | ✅ [Guide] | ✅ [Guide] | ✅ |
| Read Owners of a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Add Owner to a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Remove Owner from a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Add Lineage | ✅ [Guide] | ✅ [Guide] | ✅ |
| Add Column Level (Fine Grained) Lineage | 🚫 | ✅ [Guide] | ✅ |
| Add Documentation (Description) to a Column of a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Add Documentation (Description) to a Dataset | ✅ [Guide] | ✅ [Guide] | ✅ |
| Add / Remove / Replace Custom Properties on a Dataset | 🚫 | ✅ [Guide] | ✅ |
| Add ML Feature to ML Feature Table | 🚫 | ✅ [Guide] | ✅ |
| Add ML Feature to MLModel | 🚫 | ✅ [Guide] | ✅ |
| Add ML Group to MLFeatureTable | 🚫 | ✅ [Guide] | ✅ |
| Create MLFeature | 🚫 | ✅ [Guide] | ✅ |
| Create MLFeatureTable | 🚫 | ✅ [Guide] | ✅ |
| Create MLModel | 🚫 | ✅ [Guide] | ✅ |
| Create MLModelGroup | 🚫 | ✅ [Guide] | ✅ |
| Create MLPrimaryKey | 🚫 | ✅ [Guide] | ✅ |
| Read MLFeature | ✅ [Guide] | ✅ [Guide] | ✅ |
| Read MLFeatureTable | ✅ [Guide] | ✅ [Guide] | ✅ |
| Read MLModel | ✅ [Guide] | ✅ [Guide] | ✅ |
| Read MLModelGroup | ✅ [Guide] | ✅ [Guide] | ✅ |
| Read MLPrimaryKey | ✅ [Guide] | ✅ [Guide] | ✅ |
| Create Data Product | 🚫 | ✅ [Code] | ✅ |
| Create Lineage Between Chart and Dashboard | 🚫 | ✅ [Code] | ✅ |
| Create Lineage Between Dataset and Chart | 🚫 | ✅ [Code] | ✅ |
| Create Lineage Between Dataset and DataJob | 🚫 | ✅ [Code] | ✅ |
| Create Finegrained Lineage as DataJob for Dataset | 🚫 | ✅ [Code] | ✅ |
| Create Finegrained Lineage for Dataset | 🚫 | ✅ [Code] | ✅ |
| Create DataJob with Dataflow | 🚫 | ✅ [Code] | ✅ |
| Create Programmatic Pipeline | 🚫 | ✅ [Code] | ✅ |
注:表中 Python SDK 一列的原版链接指向 GitHub 上
metadata-ingestion/examples/library/目录下的示例脚本,本文已按当前仓库结构转换为对应的仓库相对路径,便于读者直接查看可运行的示例代码。
七、选型决策建议与深入学习路径
综合上述 API 定位、能力矩阵与底层事件模型,给出如下选型建议:
- UI 能力范围内的交互式查询与更新(搜索、关系查询、软删除、打标/术语/域/所有者等)→ 优先考虑GraphQL API,其能力与前端镜像对齐、上手直观,且天然支持 UI 场景的缓存与同步语义;
- 程序化批量写入、数据摄取、自定义实体建模、血缘注入、列级精细血缘→ 使用Python SDK / Java SDK(官方最推荐),并配合
ASYNC发射模式换取吞吐; - 需要最底层、最强灵活的写入能力(如直接 UPSERT/CREATE 任意实体+方面、批量读取、版本化回读、通用 Patch)→ 使用OpenAPI v3,其规范可自动生成客户端代码;
- 离线/解耦场景→ 使用Kafka Emitter或 JavaFile Emitter,将元数据先写入消息总线或 JSON 文件再异步导入;
- 安全红线:任何情况下都不要绕过 GMS 直接向 Kafka 写入 MCP,以免失去授权与校验保护。
继续深入的学习路径:
- Python SDK / Emitter 完整指南
- Java SDK / Emitter 完整指南 与 Java SDK V2
- GraphQL 入门指南、GraphQL 最佳实践 与 GraphQL 端点开发指南
- OpenAPI 使用指南
- MCP/MCL 事件模型详解
- API 分场景教程合集(数据集、标签、术语、域、所有者、血缘、描述、自定义属性、ML 实体等)
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考