返回博客

元数据目录:Ontology 驱动的数据资产发现与治理

coomia-dip 的元数据目录以 Ontology 为核心,整合了技术元数据(Schema、统计信息)、业务元数据(描述、标签、所有者)和操作元数据(血缘、分类、质量指标)。目录提供统一搜索、数据地图、影响分析和合规视图,支持通过 gRPC API 和 SDK 进行程序化访问。本文从目录架构、元数据模型、搜索引擎、数据地图到治理工作流,完整解析这一企业级元数据目录能力。

Coomia发布于 2025年9月25日10 分钟阅读
分享本文Twitter / X

系列:S6 平台工程 · 第 11 篇 | 难度:高级 | 阅读时间:18 分钟

元数据目录:Ontology 驱动的数据资产发现与治理

#TL;DR

coomia-dip 的元数据目录以 Ontology 为核心,整合了技术元数据(Schema、统计信息)、业务元数据(描述、标签、所有者)和操作元数据(血缘、分类、质量指标)。目录提供统一搜索、数据地图、影响分析和合规视图,支持通过 gRPC API 和 SDK 进行程序化访问。本文从目录架构、元数据模型、搜索引擎、数据地图到治理工作流,完整解析这一企业级元数据目录能力。

#1. 元数据目录的核心价值

#1.1 数据发现的挑战

随着 Ontology 对象类型增长到数百甚至数千个,用户面临"数据在哪里"的困境:

  • 发现困难:不知道平台上有哪些数据资产可用
  • 理解困难:找到数据后不理解字段含义和业务上下文
  • 信任困难:不确定数据的质量、新鲜度和可靠性
  • 合规困难:无法快速定位包含敏感数据的资产

#1.2 Ontology-First 设计

与传统元数据目录(如 Apache Atlas、DataHub)不同,coomia-dip 的目录以 Ontology 为第一公民:

Code
┌──────────────────────────────────────────────┐
│             Metadata Catalog                  │
│  ┌────────────────────────────────────────┐  │
│  │         Ontology Layer                 │  │
│  │  ObjectTypes, Links, Actions           │  │
│  │  (核心数据模型 = 元数据骨架)             │  │
│  └────────────────┬───────────────────────┘  │
│                   │                          │
│  ┌────────┐ ┌─────┴──────┐ ┌──────────────┐ │
│  │Technical│ │ Business   │ │ Operational  │ │
│  │Metadata │ │ Metadata   │ │ Metadata     │ │
│  │         │ │            │ │              │ │
│  │• Schema │ │• 描述      │ │• 血缘        │ │
│  │• 统计   │ │• 标签      │ │• 分类        │ │
│  │• 分区   │ │• 所有者    │ │• 质量指标    │ │
│  │• 格式   │ │• 领域      │ │• 使用统计    │ │
│  └────────┘ └────────────┘ └──────────────┘ │
└──────────────────────────────────────────────┘

#1.3 对标 Palantir Foundry

能力Palantir Foundrycoomia-dip
元数据核心Dataset 为中心Ontology 为中心
搜索全文搜索全文 + 语义搜索
数据地图Monocle内置拓扑图
标签系统支持层级标签 + 自动标签
API内部gRPC + SDK

#2. 元数据模型

#2.1 资产模型

Python
class CatalogAsset(BaseModel):
    """元数据目录资产"""

    asset_id: str = Field(description="资产唯一标识")
    asset_type: AssetType = Field(description="资产类型")

    # 基本信息
    name: str
    display_name: str
    description: str = ""
    namespace: str = "default"

    # 业务元数据
    owner: AssetOwner
    domain: str = Field(description="业务领域")
    tags: list[Tag] = Field(default_factory=list)
    glossary_terms: list[str] = Field(default_factory=list)

    # 技术元数据
    technical: TechnicalMetadata

    # 操作元数据
    operational: OperationalMetadata

    # 安全元数据
    classification: ClassificationLevel
    sensitivity_tags: list[str] = Field(default_factory=list)

    # 质量评分
    quality_score: float = Field(default=0.0, ge=0.0, le=1.0)
    trust_score: float = Field(default=0.0, ge=0.0, le=1.0)

    # 时间
    created_at: datetime
    updated_at: datetime
    last_accessed_at: datetime | None = None


class AssetType(str, Enum):
    OBJECT_TYPE = "object_type"
    LINK_TYPE = "link_type"
    ACTION_TYPE = "action_type"
    DATASET = "dataset"
    PIPELINE = "pipeline"
    DASHBOARD = "dashboard"


class TechnicalMetadata(BaseModel):
    """技术元数据"""
    schema_version: int
    properties: list[PropertyMetadata]
    storage_format: str
    partition_spec: dict | None = None
    row_count: int | None = None
    size_bytes: int | None = None
    last_modified: datetime | None = None


class PropertyMetadata(BaseModel):
    """属性元数据"""
    name: str
    display_name: str
    data_type: str
    nullable: bool = True
    description: str = ""
    classification: ClassificationLevel | None = None
    statistics: PropertyStatistics | None = None


class PropertyStatistics(BaseModel):
    """属性统计信息"""
    distinct_count: int | None = None
    null_count: int | None = None
    min_value: str | None = None
    max_value: str | None = None
    avg_length: float | None = None
    sample_values: list[str] = Field(default_factory=list)


class OperationalMetadata(BaseModel):
    """操作元数据"""
    lineage_available: bool = False
    upstream_count: int = 0
    downstream_count: int = 0
    pipeline_ids: list[str] = Field(default_factory=list)
    refresh_frequency: str | None = None
    sla_target: str | None = None
    last_refresh: datetime | None = None
    access_count_30d: int = 0
    unique_users_30d: int = 0

#2.2 标签体系

Python
class Tag(BaseModel):
    """元数据标签"""
    key: str
    value: str | None = None
    category: TagCategory = TagCategory.USER_DEFINED
    source: TagSource = TagSource.MANUAL
    confidence: float = 1.0

class TagCategory(str, Enum):
    DOMAIN = "domain"           # 业务领域
    SENSITIVITY = "sensitivity" # 敏感度
    LIFECYCLE = "lifecycle"     # 生命周期
    QUALITY = "quality"         # 质量
    COMPLIANCE = "compliance"   # 合规
    USER_DEFINED = "user_defined"

class TagSource(str, Enum):
    MANUAL = "manual"
    AUTO_CLASSIFIED = "auto_classified"
    INHERITED = "inherited"
    ML_SUGGESTED = "ml_suggested"

#3. 搜索引擎

#3.1 全文搜索

Python
class CatalogSearchEngine:
    """元数据目录搜索引擎"""

    async def search(
        self,
        query: str,
        filters: SearchFilters | None = None,
        sort: SortSpec | None = None,
        pagination: Pagination = Pagination(),
    ) -> SearchResult:
        """全文搜索数据资产"""
        # 构建搜索查询
        search_query = self._build_query(query)

        # 应用过滤器
        if filters:
            if filters.asset_types:
                search_query = search_query.filter_by_types(filters.asset_types)
            if filters.domains:
                search_query = search_query.filter_by_domains(filters.domains)
            if filters.classification_max:
                search_query = search_query.filter_by_classification(
                    max_level=filters.classification_max,
                )
            if filters.tags:
                search_query = search_query.filter_by_tags(filters.tags)
            if filters.owner:
                search_query = search_query.filter_by_owner(filters.owner)

        # 执行搜索
        results = await self._search_index.execute(
            search_query,
            offset=pagination.offset,
            limit=pagination.limit,
        )

        return SearchResult(
            total=results.total,
            items=[self._to_search_hit(r) for r in results.hits],
            facets=results.facets,
        )

    async def suggest(self, prefix: str, limit: int = 10) -> list[Suggestion]:
        """搜索建议(自动补全)"""
        return await self._search_index.suggest(prefix, limit)

#3.2 数据地图

Python
class DataMap:
    """数据地图 - 可视化数据资产拓扑"""

    async def get_domain_map(self, domain: str | None = None) -> DomainTopology:
        """获取领域级数据地图"""
        assets = await self._catalog.list_assets(domain=domain)

        nodes = []
        edges = []

        for asset in assets:
            nodes.append(MapNode(
                id=asset.asset_id,
                label=asset.display_name,
                type=asset.asset_type,
                classification=asset.classification,
                quality_score=asset.quality_score,
            ))

            # 添加血缘边
            lineage = await self._lineage_query.get_downstream(asset.asset_id, max_depth=1)
            for edge in lineage.edges:
                edges.append(MapEdge(
                    source=edge.source_id,
                    target=edge.target_id,
                    type=edge.edge_type,
                ))

        return DomainTopology(nodes=nodes, edges=edges, domain=domain)

    async def get_asset_neighborhood(
        self, asset_id: str, depth: int = 2,
    ) -> AssetNeighborhood:
        """获取资产的邻域图"""
        upstream = await self._lineage_query.get_upstream(asset_id, max_depth=depth)
        downstream = await self._lineage_query.get_downstream(asset_id, max_depth=depth)

        return AssetNeighborhood(
            center=asset_id,
            upstream_graph=upstream,
            downstream_graph=downstream,
        )

#4. 数据质量评分

#4.1 质量维度

Python
class DataQualityScorer:
    """数据质量评分器"""

    DIMENSIONS = {
        "completeness": 0.25,   # 完整性权重
        "accuracy": 0.20,       # 准确性权重
        "freshness": 0.20,      # 新鲜度权重
        "consistency": 0.15,    # 一致性权重
        "documentation": 0.10,  # 文档完善度权重
        "accessibility": 0.10,  # 可访问性权重
    }

    async def compute_score(self, asset: CatalogAsset) -> QualityReport:
        """计算资产的综合质量评分"""
        scores = {}

        scores["completeness"] = await self._score_completeness(asset)
        scores["accuracy"] = await self._score_accuracy(asset)
        scores["freshness"] = self._score_freshness(asset)
        scores["consistency"] = await self._score_consistency(asset)
        scores["documentation"] = self._score_documentation(asset)
        scores["accessibility"] = self._score_accessibility(asset)

        overall = sum(
            scores[dim] * weight
            for dim, weight in self.DIMENSIONS.items()
        )

        return QualityReport(
            asset_id=asset.asset_id,
            overall_score=overall,
            dimension_scores=scores,
            computed_at=datetime.utcnow(),
        )

    def _score_freshness(self, asset: CatalogAsset) -> float:
        """评估数据新鲜度"""
        if not asset.operational.last_refresh:
            return 0.0
        age = datetime.utcnow() - asset.operational.last_refresh
        if age < timedelta(hours=1):
            return 1.0
        elif age < timedelta(hours=24):
            return 0.8
        elif age < timedelta(days=7):
            return 0.5
        elif age < timedelta(days=30):
            return 0.3
        else:
            return 0.1

    def _score_documentation(self, asset: CatalogAsset) -> float:
        """评估文档完善度"""
        score = 0.0
        if asset.description:
            score += 0.3
        if asset.owner:
            score += 0.2
        documented_props = sum(
            1 for p in asset.technical.properties if p.description
        )
        total_props = len(asset.technical.properties)
        if total_props > 0:
            score += 0.5 * (documented_props / total_props)
        return score

#5. 治理工作流

#5.1 资产认证

Python
class AssetCertification:
    """数据资产认证"""

    class CertificationLevel(str, Enum):
        UNCERTIFIED = "uncertified"
        BRONZE = "bronze"      # 基本文档齐全
        SILVER = "silver"      # 质量检查通过
        GOLD = "gold"          # 完全治理合规

    CERTIFICATION_REQUIREMENTS = {
        CertificationLevel.BRONZE: [
            "has_description",
            "has_owner",
            "has_classification",
        ],
        CertificationLevel.SILVER: [
            "has_description",
            "has_owner",
            "has_classification",
            "quality_score_above_0.7",
            "has_lineage",
        ],
        CertificationLevel.GOLD: [
            "has_description",
            "has_owner",
            "has_classification",
            "quality_score_above_0.9",
            "has_lineage",
            "has_sla",
            "all_properties_documented",
            "compliance_review_passed",
        ],
    }

    async def evaluate_certification(
        self, asset: CatalogAsset,
    ) -> CertificationResult:
        """评估资产认证等级"""
        results = {}
        for level in [
            self.CertificationLevel.GOLD,
            self.CertificationLevel.SILVER,
            self.CertificationLevel.BRONZE,
        ]:
            requirements = self.CERTIFICATION_REQUIREMENTS[level]
            met = all(self._check_requirement(asset, req) for req in requirements)
            results[level] = met

        if results[self.CertificationLevel.GOLD]:
            achieved = self.CertificationLevel.GOLD
        elif results[self.CertificationLevel.SILVER]:
            achieved = self.CertificationLevel.SILVER
        elif results[self.CertificationLevel.BRONZE]:
            achieved = self.CertificationLevel.BRONZE
        else:
            achieved = self.CertificationLevel.UNCERTIFIED

        return CertificationResult(
            asset_id=asset.asset_id,
            level=achieved,
            details=results,
        )

#5.2 所有权管理

Python
class OwnershipManager:
    """数据资产所有权管理"""

    async def assign_owner(
        self, asset_id: str, owner: AssetOwner, assigned_by: str,
    ) -> None:
        asset = await self._catalog.get_asset(asset_id)
        old_owner = asset.owner

        asset.owner = owner
        await self._catalog.update_asset(asset)

        # 记录审计
        await self._audit.emit(AuditEvent(
            event_type=AuditEventType.SCHEMA_CHANGE,
            action="assign_owner",
            changes=[AuditChange(
                field="owner",
                old_value=old_owner.model_dump_json() if old_owner else None,
                new_value=owner.model_dump_json(),
                change_type="update",
            )],
        ))

    async def find_orphan_assets(self) -> list[CatalogAsset]:
        """查找无主资产"""
        all_assets = await self._catalog.list_assets()
        return [a for a in all_assets if not a.owner or a.owner.is_inactive]

#6. gRPC 服务接口

PROTOBUF
syntax = "proto3";
package onto.catalog.v1;

service CatalogService {
    rpc SearchAssets(SearchRequest) returns (SearchResponse);
    rpc GetAsset(GetAssetRequest) returns (CatalogAsset);
    rpc ListAssets(ListAssetsRequest) returns (ListAssetsResponse);

    rpc AddTag(AddTagRequest) returns (CatalogAsset);
    rpc RemoveTag(RemoveTagRequest) returns (CatalogAsset);

    rpc GetDataMap(DataMapRequest) returns (DataMapResponse);
    rpc GetAssetNeighborhood(NeighborhoodRequest) returns (NeighborhoodResponse);

    rpc GetQualityReport(QualityRequest) returns (QualityReport);
    rpc GetCertification(CertificationRequest) returns (CertificationResult);

    rpc AssignOwner(AssignOwnerRequest) returns (CatalogAsset);
    rpc FindOrphanAssets(FindOrphansRequest) returns (ListAssetsResponse);

    rpc Suggest(SuggestRequest) returns (SuggestResponse);
}

#7. 测试策略

Python
class TestMetadataCatalog:
    async def test_search_by_keyword(self):
        results = await search_engine.search("employee salary")
        assert len(results.items) > 0
        assert any("Employee" in r.name for r in results.items)

    async def test_search_with_filters(self):
        results = await search_engine.search(
            "customer",
            filters=SearchFilters(
                asset_types=[AssetType.OBJECT_TYPE],
                classification_max=ClassificationLevel.CONFIDENTIAL,
            ),
        )
        for item in results.items:
            assert item.classification <= ClassificationLevel.CONFIDENTIAL

    async def test_quality_scoring(self):
        scorer = DataQualityScorer()
        report = await scorer.compute_score(well_documented_asset)
        assert report.overall_score > 0.7
        assert report.dimension_scores["documentation"] > 0.8

    async def test_certification_levels(self):
        cert = AssetCertification()
        result = await cert.evaluate_certification(gold_asset)
        assert result.level == AssetCertification.CertificationLevel.GOLD

    async def test_find_orphan_assets(self):
        manager = OwnershipManager(catalog)
        orphans = await manager.find_orphan_assets()
        for orphan in orphans:
            assert not orphan.owner or orphan.owner.is_inactive

#8. 生产最佳实践

#8.1 元数据采集

  • Schema 变更自动触发目录更新
  • 统计信息每天定期采集
  • 使用统计实时流式采集
  • 血缘信息从 Pipeline 执行中自动提取

#8.2 治理流程

  1. 新资产必须在 7 天内达到 Bronze 认证
  2. 面向业务用户的资产应达到 Silver 认证
  3. 关键业务资产必须达到 Gold 认证
  4. 无主资产每月检查,超过 30 天未认领则降级

#8.3 搜索优化

  • 建立领域特定的同义词表
  • 利用使用统计提升热门资产的搜索排名
  • 高质量评分的资产在搜索结果中优先展示

#9. 总结

coomia-dip 的元数据目录以 Ontology 为核心,实现了数据资产的统一发现、理解和治理。关键设计亮点:

  1. Ontology-First:以对象类型为核心组织元数据,而非传统的表/文件
  2. 三层元数据:技术、业务、操作三维元数据全面覆盖
  3. 质量评分:6 维度量化评分体系,客观评估数据可信度
  4. 认证体系:Bronze/Silver/Gold 三级认证驱动数据治理
  5. 搜索引擎:全文搜索 + 多维过滤 + 自动补全

下一篇将深入探讨 coomia-dip 的合规设计体系。