OpenMetadata Python SDK 正确 API 形态与常见迁移纠错:复数 Facade、configure() 与 CRUD、分页、Patch、搜索、血缘、CSV 的源码级详解
【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata
OpenMetadata 的 Python SDK(metadata.sdk)是当前仓库中面向外部脚本、自动化任务和 AI Agent 的高层编程入口,它把 ingestion 包内生成的 Pydantic 实体模型封装为一批"复数命名"的 Facade 类。本文基于仓库内 SDK 兼容性说明文档 的完整内容展开,并逐条对照 sdk 入口、客户端封装、基类实现 等源码,帮助读者掌握:SDK 的正确导入方式与一次性配置方法、实体 CRUD 与分页的完整用法、patch 更新的真实调用链,以及搜索、血缘、治理标签、CSV 四类操作的边界与迁移纠错点——避免把早期设计草案中的 API(如单数实体类、patch(entity_id, json_patch)、TableListParams等)误写进自己的代码。
适用前提与安装
SDK 位于ingestion子项目内,随openmetadata-ingestion包一起发布。按 SDK 主 README 的说明:
pip install openmetadata-ingestion需要数据质量等附加能力时安装对应 extras(如openmetadata-ingestion[pandas]、openmetadata-ingestion[mysql])。SDK 运行依赖一个可达的 OpenMetadata Server(默认示例地址http://localhost:8585/api)和一个 JWT 令牌。
复数 Facade 类:实体操作的唯一入口
这是整份文档最关键的一条迁移纠错:SDK 操作全部挂在复数 Facade 类上(Tables、Databases、Users、Dashboards……),而metadata.generated.schema.entity.data.table.Table这类单数生成实体类只是 Pydantic 模型,本身不暴露任何 SDK 方法。正确写法是"用 Facade 取数、用单数模型做类型标注":
from metadata.sdk import Databases, Tables, Users from metadata.generated.schema.entity.data.table import Table table: Table = Tables.retrieve_by_name("service.database.schema.orders")文档明确警告:不要从metadata.sdk.entities导入Table、Database、Dashboard等单数 Facade 名,它们不是已导出的 SDK 类。两种合法的导入方式:
from metadata.sdk import Dashboards, Databases, Tables或模块级精确导入:
from metadata.sdk.entities.tables import Tables从源码结构看,这一命名约定在 entities 包初始化 中有注释直接说明:"Plural naming convention to avoid conflicts with generated entities"(复数命名以避免与生成的实体类冲突)。所有 Facade 类都继承自 BaseEntity,它声明了统一的entity_type()契约和 CRUD 方法集,因此Tables、Glossaries、Dashboards等类的方法签名高度一致,学会一个即会全部。
一次性配置:configure() 与手动初始化
绝大多数脚本应使用configure(),它在 sdk 入口 中实现,作用是建立全局默认客户端,供所有 Facade 类共享:
from metadata.sdk import configure configure(host="http://localhost:8585/api", jwt_token="your-jwt-token")configure()支持三种传参形式(见其源码 docstring):
- 关键字参数:
host/server_url(两者等价,host是别名)加jwt_token; - 环境变量:
configure()不带参数时调用OpenMetadataConfig.from_env(),读取以下变量(定义见 config.py 的 from_env):
| 环境变量 | 作用 | 默认值 |
|---|---|---|
OPENMETADATA_HOST或OPENMETADATA_SERVER_URL | 服务端 URL(必填) | 无 |
OPENMETADATA_JWT_TOKEN或OPENMETADATA_API_KEY | 认证令牌 | 无 |
OPENMETADATA_VERIFY_SSL | 是否校验 SSL | false |
OPENMETADATA_CA_BUNDLE | CA 证书包路径 | 无 |
OPENMETADATA_CLIENT_TIMEOUT | 客户端超时(秒) | 30 |
- 显式配置对象或 mapping:注意源码中有保护性检查——同时传入
config对象和关键字参数会抛出TypeError("Pass either a config object or keyword arguments, not both")。
进阶场景(测试、多实例)可手动初始化:
from metadata.sdk import OpenMetadata, OpenMetadataConfig config = OpenMetadataConfig( server_url="http://localhost:8585/api", jwt_token="your-jwt-token", ) client = OpenMetadata.initialize(config)两条路径最终都会执行 OpenMetadata.initialize:把OpenMetadataConfig转成 ingestion 客户端内部的OpenMetadataConnection(AuthProvider=openmetadata,按verify_ssl/ca_bundle推导verifySSL枚举),再实例化底层的 OMeta 客户端。OpenMetadataConfig 的完整字段为:server_url(自动去掉末尾/)、jwt_token(jwt_token or api_key)、api_key、verify_ssl(默认False)、ca_bundle、client_timeout(默认30秒);此外还提供一个链式OpenMetadataConfig.builder()构建器。
配套的生命周期管理函数(同样在 sdk 入口 中):
from metadata.sdk import client, reset metadata = client().ometa # 获取底层客户端;未配置时抛出 RuntimeError reset() # 关闭并清空全局客户端client()在未调用configure()/initialize()前会抛出RuntimeError("SDK not configured..."),这是排查"SDK 未配置"类报错的第一现场。
CRUD 完整走查:以 Tables 为例
文档给出的 CRUD 示例完整覆盖了create → retrieve → retrieve_by_name → update → delete全链路,可直接复制运行:
from metadata.generated.schema.api.data.createTable import CreateTableRequest from metadata.generated.schema.entity.data.table import Column, DataType from metadata.generated.schema.type.basic import Markdown from metadata.sdk import Tables request = CreateTableRequest( name="orders", databaseSchema="service.database.schema", columns=[ Column(name="id", dataType=DataType.BIGINT), Column(name="status", dataType=DataType.VARCHAR, dataLength=255), ], ) table = Tables.create(request) table = Tables.retrieve(str(table.id.root), fields=["owners", "tags"]) table = Tables.retrieve_by_name("service.database.schema.orders") updated_table = table.model_copy(deep=True) updated_table.description = Markdown("Curated order facts") table = Tables.update(updated_table) Tables.delete(str(table.id.root), recursive=True, hard_delete=True)几个值得注意的源码细节(均在 base.py 中可验证):
table.id是包装类型:id字段是带root属性的包装对象,取字符串值要写str(table.id.root)。基类内部的_stringify_identifier()会自动完成这一解包,所以传table.id或字符串给 Facade 方法都可以。update()只接收实体对象:其实现(base.py#L242-L254)分三步——从entity.id读取标识、get_by_id拉取当前线上版本、再用client.patch(entity=..., source=current, destination=entity)对变更字段做增量补丁。因此推荐的修改方式是先model_copy(deep=True)再改字段,而不是手工构造请求体。delete()支持级联与硬删:delete(entity_id, recursive=False, hard_delete=False)对应 base.py#L256-L270,recursive=True用于删除表时连带删除其列等子实体。retrieve()/retrieve_by_name()支持fields与nullable:fields=["owners", "tags"]做字段级裁剪以降低传输量;nullable=True时实体不存在返回None而非抛异常。
TablesFacade 还有一批表专属便捷方法(见 tables.py),文档之外的实用能力包括:update_column_description(table_id, column_name, description)只更新某一列的描述、add_custom_metric()、add_sample_data()/get_sample_data()管理样例数据。
列表与分页:list() 与 list_all()
list()返回的是单页EntityList对象(base.py#L25-L31 中的 dataclass,含entities、after、before三个字段),需要手动用游标翻页:
from metadata.sdk import Tables page = Tables.list( limit=50, fields=["owners", "tags"], filters={"databaseSchema": "service.database.schema"}, ) for table in page.entities: print(table.name) if page.after: next_page = Tables.list(limit=50, after=page.after)需要遍历全部实体时,用list_all()自动翻页(内部循环调用list(),直到after为空,实现见 base.py#L300-L323):
for table in Tables.list_all(batch_size=100): print(table.fullyQualifiedName)两条重要的迁移纠错:
- 不存在
TableListParams类——早期设计草案中的分页参数类已废弃,参数直接以关键字形式传给list()/list_all(); EntityList没有auto_paging_iterable()方法——自动翻页只通过list_all()提供。
从源码看默认值:list()的limit默认是10,list_all()的batch_size默认是100,filters是一个普通dict[str, str],会被原样透传为查询参数。
Patch 行为:没有 patch(entity_id, json_patch)
Facade 层不暴露patch(entity_id, json_patch)这种方法。普通的局部更新一律走update(entity)(如上节所述,内部自动完成"取旧值 → diff → 补丁")。若确实需要低层 patch 流程(例如精细控制 JSON Patch),可以下探到默认客户端:
from metadata.generated.schema.entity.data.table import Table from metadata.sdk import client metadata = client().ometa source = metadata.get_by_id(entity=Table, entity_id="table-id", fields=["tags"]) destination = source.model_copy(deep=True) destination.tags = [] patched = metadata.patch(entity=Table, source=source, destination=destination)这个 source/destination 双实体模式正是 Facade 内部update()的实现方式——把 base.py 的 update() 实现 摊开写而已。换句话说,Tables.update()与这里的低层调用走的是同一条metadata.patch通道,只是前者帮你省去了手动取 source 的步骤。
异步支持的边界:搜索与血缘可用,实体 CRUD 仅同步
文档对异步能力划了一条清晰的线:Search和Lineage暴露*_async()辅助方法;实体 CRUD Facade 目前只有同步方法,示例代码不应出现Tables.create_async()、Tables.retrieve_async()之类不存在的调用。
from metadata.sdk.api import Search results = await Search.search_async("customer", index="table_search_index")从源码看(search.py),这些 async 方法是_run_async()对同步方法的包装,底层通过loop.run_in_executor(None, func)把阻塞调用丢到默认线程池执行——它解决的是"不阻塞事件循环",而不是真正的网络层异步。Search可用的完整方法集包括search(支持from_/size/sort_field/sort_order/filters/include_aggregations)、suggest、aggregate、search_advanced、reindex/reindex_all,以及对应的SearchBuilder链式构建器;Lineage同理提供get_lineage_async、add_lineage_async等(lineage.py)。
治理标签:Classifications 与 Tags
治理分类体系(taxonomy)的增删改查使用Classifications和Tags两个 Facade;给表打标签有两条路径:
- 追加一个标签:
Tables.add_tag(table_id, "Classification.Tag"); - 整体重排标签:先以
fields=["tags"]取表,在副本上替换tags,再调Tables.update(destination)。
from metadata.generated.schema.api.classification.createClassification import ( CreateClassificationRequest, ) from metadata.generated.schema.api.classification.createTag import CreateTagRequest from metadata.sdk import Classifications, Tables, Tags classification = Classifications.create( CreateClassificationRequest(name="PII", description="PII taxonomy") ) tag = Tags.create( CreateTagRequest( classification=classification.fullyQualifiedName.root, name="Sensitive", description="Sensitive data", ), ) table = Tables.add_tag("table-id", tag.fullyQualifiedName.root)Tables.add_tag的源码实现(tables.py#L29-L51)展示了它为什么是幂等追加:先get_by_id(..., fields=["tags"])只取 tags 字段,深拷贝后tags.append({"tagFQN": tag_fqn}),最后patch(source=current, destination=working)——即每次调用只会让标签集合新增一项,不会误删已有标签。需要批量重排时,按第二条路径用update()全量替换即可。
血缘:在 metadata.sdk.api.Lineage 下
血缘操作不在实体 Facade 下,而是位于metadata.sdk.api.Lineage。两个核心用法:按实体 ID 加边,或按 FQN/ID 取血缘图:
from metadata.generated.schema.entity.data.table import Table from metadata.sdk.api import Lineage Lineage.add_lineage( from_entity_id="source-table-id", from_entity_type="table", to_entity_id="target-table-id", to_entity_type="table", description="Curated order facts", ) lineage = Lineage.get_entity_lineage( entity_type=Table, entity_id="target-table-id", upstream_depth=2, downstream_depth=1, )源码中(lineage.py)还能看到文档示例之外的实用方法:add_lineage_by_name()直接用两端 FQN 建边、get_lineage()按 FQN 查血缘、delete_lineage()/delete_lineage_by_name()删边、export_lineage()把血缘图导出为 JSON(默认上下游各 3 层),以及链式LineageBuilder。add_lineage的实现会先把 ID 做ensure_uuid校验、把description包装成Markdown模型,再构造AddLineageRequest(edge=EntitiesEdge(...))提交给底层客户端——这也是理解血缘请求体结构的最短路径。
CSV 操作:返回操作对象,再 execute()
CSV 导出/导入是流式操作对象模式:方法调用先拿到一个 Operation 对象,配置完再execute():
from metadata.sdk import Glossaries, Tables csv_text = Tables.export_csv("service.database.schema.orders").execute() dry_run = ( Glossaries.import_csv("BusinessGlossary") .with_data(csv_text) .set_dry_run(True) .execute() )源码对应 base.py 中的 CsvExportOperation / CsvImportOperation:CsvImportOperation持有csv_data与dry_run两个状态(set_dry_run(True)对应服务端"只校验不落库"的试运行),execute()内部调用client.import_csv(entity=..., name=..., csv_data=..., dry_run=...);CsvExportOperation.execute()内部调用client.export_csv(...)。两者都还有execute_async()变体(若底层客户端不支持则抛AttributeError)。
当前 Facade 覆盖范围与缺口的替代方案
文档列出的当前导出清单(可直接复制做导入自检):
from metadata.sdk import ( APICollections, APIEndpoints, Charts, Classifications, Containers, DashboardDataModels, DashboardServices, Dashboards, DatabaseSchemas, DatabaseServices, Databases, DataContracts, DataProducts, Domains, Glossaries, GlossaryTerms, Metrics, MLModels, Pipelines, Queries, SearchIndexes, StorageServices, StoredProcedures, Tables, Tags, Teams, TestCases, TestDefinitions, TestSuites, Users, )对照当前 sdk 入口的__all__,源码实际上还额外导出了ContextFiles、Folders、Pages、Settings四个 Facade 及BaseEntity、configure、client、reset、to_entity_reference等工具符号,并提供小写别名(如tables、glossary_terms)方便按stripe风格使用——以源码清单为准更稳妥。
对于还没有 Facade 的实体,文档给出的替代方案是直接使用metadata.sdk.client()返回的底层客户端,或回到 ingestion 的 OpenMetadata 客户端(client().ometa)做通用操作,直到对应 Facade 补齐。另外,to_entity_reference(entity)是设置owners、domains等引用字段时的标准工具:它把任意实体转成只含id/type/name/fullyQualifiedName的 EntityReference dict(入口实现)。
验证与延伸阅读
上述 Facade 行为在仓库中有成体系的单元测试可作对照:test_base_entity.py 覆盖基类 CRUD 与分页,test_table_entity.py 覆盖add_tag等表专属方法,test_csv_operations.py、test_sdk_apis.py(搜索/血缘 API)、test_restore_async.py(异步恢复的 202 job 响应形态)分别对应本文各节的 API。
最后归纳本文涉及的迁移纠错清单,方便自查存量代码:
| 旧草案写法(不要再用) | 当前正确写法 |
|---|---|
从metadata.sdk.entities导入单数类Table/Database/Dashboard当 Facade 用 | 导入复数Tables/Databases/Dashboards;单数名仅作 Pydantic 模型 |
Tables.patch(table_id, json_patch) | Tables.update(entity);低层场景走client().ometa.patch(...) |
TableListParams参数类、EntityList.auto_paging_iterable() | list(limit=..., after=..., filters=...)与list_all(batch_size=...) |
Tables.create_async()/Tables.retrieve_async() | 实体 CRUD 目前仅同步;异步仅Search/Lineage的*_async() |
期望 Facade 上有patch(entity_id, ...)风格的浅更新 | model_copy(deep=True)改字段后update() |
掌握以上形态后,编写脚本时的标准范式是:configure()一次 → 复数 Facade 做 CRUD/列表 →Search/Lineage做检索与血缘 → 无 Facade 的实体下探client().ometa,所有示例均可在 README_IMPROVED.md 与 sdk 源码目录 中找到逐行可验证的依据。
【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考