返回博客

entity_common / entity_edge / entity_event:三表模型设计

Tags: #ThreeTableModel #Ontology #SchemaDesign #QueryPatterns #EntityModel #智策平台

Coomia发布于 2025年7月14日21 分钟阅读
分享本文Twitter / X

系列:S3 数据基座 · 第 5 篇 | 难度:高级 | 阅读时间:20 分钟

entity_common / entity_edge / entity_event:三表模型设计

Tags: #ThreeTableModel #Ontology #SchemaDesign #QueryPatterns #EntityModel #智策平台

#TL;DR

在 coomia-dip 数据基座中,我们用三张核心表——entity_common(实体表)、entity_edge(关系边表)、entity_event(事件表)——承载任意 Ontology 的全部数据。这种"三表模型"是连接上层 Ontology 语义和底层存储引擎的关键桥梁。本文详细解析为什么选择三表而非宽表或图数据库,表结构设计的每一个字段决策,以及在 Doris 和 Iceberg 上的查询优化策略。通过真实的性能对比数据,证明三表模型在灵活性和性能之间达到了最佳平衡。

#1. 为什么是三张表

#1.1 Ontology 数据的本质分类

Ontology 本体论中的数据可以归纳为三种基本类型:

Code
Ontology 数据三元分类:

┌─────────────────────────────────────────────────────────┐
│                    Ontology 数据                         │
│                                                          │
│  ┌───────────────┐ ┌───────────────┐ ┌───────────────┐  │
│  │    Entity      │ │     Edge      │ │    Event      │  │
│  │    (实体)      │ │    (关系)     │ │    (事件)     │  │
│  │               │ │               │ │               │  │
│  │ "什么东西"     │ │ "谁和谁有关"  │ │ "发生了什么"  │  │
│  │               │ │               │ │               │  │
│  │ Person        │ │ WorksAt       │ │ Login         │  │
│  │ Company       │ │ Manages       │ │ Transaction   │  │
│  │ Device        │ │ ConnectedTo   │ │ Alert         │  │
│  │ Document      │ │ DerivedFrom   │ │ StatusChange  │  │
│  └───────┬───────┘ └───────┬───────┘ └───────┬───────┘  │
│          │                 │                 │           │
│          ▼                 ▼                 ▼           │
│   entity_common      entity_edge      entity_event      │
│                                                          │
└─────────────────────────────────────────────────────────┘

#1.2 其他方案的问题

在确定三表模型之前,我们评估了多种方案:

方案描述优点致命缺陷
每类型一表Person 表、Company 表...查询简单类型爆炸,DDL 不可控
单宽表所有数据一张表极简稀疏列浪费,Schema 膨胀
EAV 模型Entity-Attribute-Value完全灵活查询性能灾难
图数据库Neo4j / JanusGraph图遍历快聚合分析弱,运维复杂
三表模型实体 + 关系 + 事件平衡适度复杂

#1.3 三表模型的设计哲学

Code
三表模型设计原则:

原则 1:类型无关的表结构
  所有 Entity 类型共享同一张表,类型区分靠 entity_type 字段
  → 新增 Ontology 类型不需要 DDL 变更

原则 2:属性 JSON 化
  每个实体的特有属性存储在 JSON 列中
  → 灵活性 > 查询性能(通过索引弥补)

原则 3:关系显式建模
  关系不是"外键",而是一等公民,拥有自己的属性
  → 支持带权重、带时间的关系

原则 4:事件流与状态分离
  实体表存"当前状态",事件表存"变化历史"
  → 天然支持事件溯源和时间旅行

#2. entity_common 表:实体的统一载体

#2.1 表结构设计

SQL
-- entity_common:所有 Ontology 实体的统一存储
-- 存储引擎:Apache Doris (主查询) + Apache Iceberg (版本化)

CREATE TABLE entity_common (
    -- === 主键区 ===
    entity_id       VARCHAR(64)     NOT NULL    COMMENT '实体唯一标识,UUID v7(时间有序)',
    world_id        VARCHAR(32)     NOT NULL    COMMENT '所属 World(项目空间)',

    -- === 类型区 ===
    entity_type     VARCHAR(128)    NOT NULL    COMMENT 'Ontology 类型,如 Person, Company',
    type_version    INT             NOT NULL    COMMENT '类型版本号,Schema 演化用',

    -- === 展示区 ===
    display_name    VARCHAR(512)    DEFAULT ''  COMMENT '显示名称(冗余,加速搜索)',
    description     TEXT            DEFAULT ''  COMMENT '描述(冗余,加速全文检索)',

    -- === 属性区 ===
    properties      JSON            NOT NULL    COMMENT '实体属性(Ontology 定义的字段)',
    computed_props  JSON            DEFAULT '{}' COMMENT '计算属性(派生属性缓存)',

    -- === 安全区 ===
    security_label  VARCHAR(32)     DEFAULT 'internal'  COMMENT '安全标签',
    owner_id        VARCHAR(64)     DEFAULT ''  COMMENT '所有者用户 ID',

    -- === 审计区 ===
    created_at      DATETIME        NOT NULL    COMMENT '创建时间',
    updated_at      DATETIME        NOT NULL    COMMENT '最后更新时间',
    created_by      VARCHAR(64)     DEFAULT ''  COMMENT '创建者',
    updated_by      VARCHAR(64)     DEFAULT ''  COMMENT '更新者',
    version         BIGINT          DEFAULT 1   COMMENT '乐观锁版本号',
    is_deleted      BOOLEAN         DEFAULT FALSE COMMENT '软删除标记',

    -- === 向量区(可选)===
    embedding       ARRAY<FLOAT>    NULL        COMMENT '语义向量(HNSW 索引)'
)
UNIQUE KEY (entity_id, world_id)
DISTRIBUTED BY HASH(entity_id) BUCKETS 32
PROPERTIES (
    "replication_num" = "3",
    "enable_unique_key_merge_on_write" = "true",
    "store_row_column" = "true"
);

-- === 索引 ===
-- 类型查询加速
CREATE INDEX idx_entity_type ON entity_common (entity_type) USING BITMAP;
-- 全文检索
CREATE INDEX idx_display_name ON entity_common (display_name) USING INVERTED
    PROPERTIES("parser" = "unicode", "support_phrase" = "true");
CREATE INDEX idx_description ON entity_common (description) USING INVERTED
    PROPERTIES("parser" = "unicode", "support_phrase" = "true");
-- JSON 属性索引
CREATE INDEX idx_properties ON entity_common (properties) USING INVERTED;
-- 向量搜索
CREATE INDEX idx_embedding ON entity_common (embedding) USING INVERTED
    PROPERTIES("index_type" = "hnsw", "metric_type" = "cosine", "m" = "16", "ef" = "200");
-- 时间范围查询
CREATE INDEX idx_created_at ON entity_common (created_at) USING INVERTED;
CREATE INDEX idx_updated_at ON entity_common (updated_at) USING INVERTED;

#2.2 字段设计决策解析

Code
为什么每个字段这样设计:

┌──────────────┬─────────────────────────────────────────────┐
│ 字段          │ 设计决策                                     │
├──────────────┼─────────────────────────────────────────────┤
│ entity_id    │ UUID v7 而非自增 ID:                         │
│              │  - 分布式环境无需协调                          │
│              │  - 时间有序,适合 B+ 树索引                    │
│              │  - 64 字符 VARCHAR 而非 BINARY:调试友好        │
├──────────────┼─────────────────────────────────────────────┤
│ world_id     │ 复合主键的一部分:                             │
│              │  - 强制数据隔离                                │
│              │  - 查询自动限定 World 范围                      │
│              │  - 支持多租户和分支                             │
├──────────────┼─────────────────────────────────────────────┤
│ entity_type  │ VARCHAR(128) 而非 ENUM:                      │
│              │  - 新增类型无需 ALTER TABLE                    │
│              │  - BITMAP 索引补偿查询性能                     │
│              │  - 存储 Ontology 完全限定名                    │
├──────────────┼─────────────────────────────────────────────┤
│ properties   │ JSON 而非列式存储:                            │
│              │  - 不同类型有不同属性,列式不可行                │
│              │  - Doris JSON 索引提供查询加速                  │
│              │  - 类型检查在应用层(SDK)完成                  │
├──────────────┼─────────────────────────────────────────────┤
│ display_name │ 冗余字段(可从 properties 导出):             │
│              │  - 避免 JSON 解析的搜索开销                    │
│              │  - 全文检索需要独立列                          │
│              │  - 列表页展示不需要解析整个 JSON                │
├──────────────┼─────────────────────────────────────────────┤
│ embedding    │ ARRAY<FLOAT> 可选列:                         │
│              │  - 仅语义搜索场景需要                          │
│              │  - 768/1536 维度视模型而定                     │
│              │  - HNSW 索引加速相似度计算                     │
└──────────────┴─────────────────────────────────────────────┘

#2.3 Properties JSON 的内部结构约定

JSON
{
  "__type_meta": {
    "type": "Person",
    "version": 3,
    "schema_hash": "a1b2c3d4"
  },
  "name": "张三",
  "email": "zhangsan@example.com",
  "age": 35,
  "department": "Engineering",
  "hire_date": "2020-06-15",
  "skills": ["Python", "Java", "gRPC"],
  "address": {
    "city": "Shanghai",
    "district": "Pudong"
  },
  "__refs": {
    "avatar": "s3://onto-proj001-uploads/avatars/zhangsan.jpg",
    "resume": "s3://onto-proj001-uploads/docs/zhangsan-resume.pdf"
  }
}

#3. entity_edge 表:关系的一等公民

#3.1 表结构设计

SQL
-- entity_edge:所有 Ontology 关系的统一存储
CREATE TABLE entity_edge (
    -- === 主键区 ===
    edge_id         VARCHAR(64)     NOT NULL    COMMENT '边唯一标识',
    world_id        VARCHAR(32)     NOT NULL    COMMENT '所属 World',

    -- === 端点区 ===
    source_id       VARCHAR(64)     NOT NULL    COMMENT '源实体 ID',
    source_type     VARCHAR(128)    NOT NULL    COMMENT '源实体类型',
    target_id       VARCHAR(64)     NOT NULL    COMMENT '目标实体 ID',
    target_type     VARCHAR(128)    NOT NULL    COMMENT '目标实体类型',

    -- === 关系区 ===
    edge_type       VARCHAR(128)    NOT NULL    COMMENT '关系类型,如 WorksAt, Manages',
    edge_label      VARCHAR(256)    DEFAULT ''  COMMENT '关系标签(显示用)',
    direction       VARCHAR(16)     DEFAULT 'directed' COMMENT 'directed/undirected',

    -- === 属性区 ===
    properties      JSON            DEFAULT '{}' COMMENT '关系属性',
    weight          DOUBLE          DEFAULT 1.0 COMMENT '关系权重',
    confidence      DOUBLE          DEFAULT 1.0 COMMENT '关系置信度',

    -- === 时间区 ===
    valid_from      DATETIME        NULL        COMMENT '关系生效时间',
    valid_to        DATETIME        NULL        COMMENT '关系失效时间',
    created_at      DATETIME        NOT NULL    COMMENT '创建时间',
    updated_at      DATETIME        NOT NULL    COMMENT '最后更新时间',

    -- === 审计区 ===
    created_by      VARCHAR(64)     DEFAULT ''  COMMENT '创建者',
    is_deleted      BOOLEAN         DEFAULT FALSE COMMENT '软删除标记',
    version         BIGINT          DEFAULT 1   COMMENT '乐观锁版本号'
)
UNIQUE KEY (edge_id, world_id)
DISTRIBUTED BY HASH(edge_id) BUCKETS 16
PROPERTIES (
    "replication_num" = "3",
    "enable_unique_key_merge_on_write" = "true"
);

-- === 索引 ===
-- 正向查询:从源实体找目标
CREATE INDEX idx_source ON entity_edge (source_id) USING BITMAP;
-- 反向查询:从目标实体找源
CREATE INDEX idx_target ON entity_edge (target_id) USING BITMAP;
-- 关系类型过滤
CREATE INDEX idx_edge_type ON entity_edge (edge_type) USING BITMAP;
-- 复合查询:特定类型的源实体的特定关系
CREATE INDEX idx_source_type ON entity_edge (source_type) USING BITMAP;
CREATE INDEX idx_target_type ON entity_edge (target_type) USING BITMAP;
-- 时间范围查询
CREATE INDEX idx_valid_from ON entity_edge (valid_from) USING INVERTED;
CREATE INDEX idx_valid_to ON entity_edge (valid_to) USING INVERTED;

#3.2 关系类型与方向

Code
关系的方向性:

1. 有向关系(Directed):
   Person ──WorksAt──> Company
   Person ──Manages──> Person
   Document ──DerivedFrom──> Document

2. 无向关系(Undirected):
   Person ──Colleagues──  Person
   Device ──ConnectedTo── Device

3. 时态关系(Temporal):
   Person ──WorksAt──> Company
     valid_from: 2020-06-15
     valid_to:   2024-03-01  (已离职)

   Person ──WorksAt──> Company
     valid_from: 2024-03-15
     valid_to:   NULL        (当前在职)

关系属性示例:
┌────────────────┬──────────────────────────┐
│ edge_type      │ properties 示例           │
├────────────────┼──────────────────────────┤
│ WorksAt        │ {"role": "Engineer",      │
│                │  "level": "Senior"}       │
├────────────────┼──────────────────────────┤
│ Manages        │ {"team_size": 8,          │
│                │  "since": "2023-01"}      │
├────────────────┼──────────────────────────┤
│ Transaction    │ {"amount": 50000,         │
│                │  "currency": "CNY"}       │
├────────────────┼──────────────────────────┤
│ Similarity     │ {"score": 0.87,           │
│                │  "method": "cosine"}      │
└────────────────┴──────────────────────────┘

#3.3 图遍历查询模式

SQL
-- 1. 一跳查询:找某人所在的公司
SELECT ec.*
FROM entity_common ec
JOIN entity_edge ee ON ee.target_id = ec.entity_id
    AND ee.world_id = ec.world_id
WHERE ee.source_id = 'person-001'
    AND ee.edge_type = 'WorksAt'
    AND ee.world_id = 'world-main'
    AND ee.is_deleted = FALSE
    AND (ee.valid_to IS NULL OR ee.valid_to > NOW());

-- 2. 两跳查询:找某人同事的项目
SELECT DISTINCT ec_project.*
FROM entity_edge ee1
JOIN entity_edge ee2 ON ee2.source_id = ee1.target_id
    AND ee2.world_id = ee1.world_id
JOIN entity_common ec_project ON ec_project.entity_id = ee2.target_id
    AND ec_project.world_id = ee2.world_id
WHERE ee1.source_id = 'person-001'
    AND ee1.edge_type = 'Colleagues'
    AND ee2.edge_type = 'WorksOn'
    AND ee1.world_id = 'world-main';

-- 3. 聚合查询:每个部门的人数
SELECT
    JSON_EXTRACT(ec.properties, '$.department') AS department,
    COUNT(*) AS headcount
FROM entity_common ec
WHERE ec.entity_type = 'Person'
    AND ec.world_id = 'world-main'
    AND ec.is_deleted = FALSE
GROUP BY JSON_EXTRACT(ec.properties, '$.department')
ORDER BY headcount DESC;

-- 4. 路径查询(递归 CTE):找两个实体之间的最短路径
WITH RECURSIVE paths AS (
    -- 起点
    SELECT
        source_id,
        target_id,
        edge_type,
        1 AS depth,
        ARRAY[source_id, target_id] AS path
    FROM entity_edge
    WHERE source_id = 'person-001'
        AND world_id = 'world-main'
        AND is_deleted = FALSE

    UNION ALL

    -- 递归展开
    SELECT
        p.source_id,
        ee.target_id,
        ee.edge_type,
        p.depth + 1,
        ARRAY_APPEND(p.path, ee.target_id)
    FROM paths p
    JOIN entity_edge ee ON ee.source_id = p.target_id
        AND ee.world_id = 'world-main'
        AND ee.is_deleted = FALSE
    WHERE p.depth < 5
        AND NOT ARRAY_CONTAINS(p.path, ee.target_id)
)
SELECT * FROM paths
WHERE target_id = 'person-099'
ORDER BY depth ASC
LIMIT 1;

#4. entity_event 表:时间线的忠实记录

#4.1 表结构设计

SQL
-- entity_event:所有事件的时间序列存储
CREATE TABLE entity_event (
    -- === 主键区 ===
    event_id        VARCHAR(64)     NOT NULL    COMMENT '事件唯一标识',
    world_id        VARCHAR(32)     NOT NULL    COMMENT '所属 World',

    -- === 关联区 ===
    entity_id       VARCHAR(64)     NOT NULL    COMMENT '关联实体 ID',
    entity_type     VARCHAR(128)    NOT NULL    COMMENT '关联实体类型',

    -- === 事件区 ===
    event_type      VARCHAR(128)    NOT NULL    COMMENT '事件类型,如 Created, Updated, Alert',
    event_subtype   VARCHAR(128)    DEFAULT ''  COMMENT '事件子类型',
    event_source    VARCHAR(128)    NOT NULL    COMMENT '事件来源',

    -- === 数据区 ===
    payload         JSON            NOT NULL    COMMENT '事件负载(完整数据)',
    summary         VARCHAR(1024)   DEFAULT ''  COMMENT '事件摘要(快速展示)',
    severity        VARCHAR(16)     DEFAULT 'info' COMMENT '严重程度: info/warn/error/critical',

    -- === 时间区 ===
    event_time      DATETIME        NOT NULL    COMMENT '事件发生时间(业务时间)',
    ingestion_time  DATETIME        NOT NULL    COMMENT '事件入库时间(系统时间)',

    -- === 溯源区 ===
    caused_by       VARCHAR(64)     DEFAULT ''  COMMENT '触发者(用户/系统/Pipeline)',
    correlation_id  VARCHAR(64)     DEFAULT ''  COMMENT '关联追踪 ID',

    -- === 审计区 ===
    is_deleted      BOOLEAN         DEFAULT FALSE COMMENT '软删除标记'
)
DUPLICATE KEY (event_id, world_id, event_time)
PARTITION BY RANGE(event_time) ()
DISTRIBUTED BY HASH(entity_id) BUCKETS 32
PROPERTIES (
    "replication_num" = "3",
    "dynamic_partition.enable" = "true",
    "dynamic_partition.time_unit" = "DAY",
    "dynamic_partition.start" = "-90",
    "dynamic_partition.end" = "3",
    "dynamic_partition.prefix" = "p",
    "dynamic_partition.buckets" = "32"
);

-- === 索引 ===
CREATE INDEX idx_entity_id ON entity_event (entity_id) USING BITMAP;
CREATE INDEX idx_entity_type ON entity_event (entity_type) USING BITMAP;
CREATE INDEX idx_event_type ON entity_event (event_type) USING BITMAP;
CREATE INDEX idx_severity ON entity_event (severity) USING BITMAP;
CREATE INDEX idx_event_source ON entity_event (event_source) USING BITMAP;
CREATE INDEX idx_summary ON entity_event (summary) USING INVERTED
    PROPERTIES("parser" = "unicode");
CREATE INDEX idx_payload ON entity_event (payload) USING INVERTED;
CREATE INDEX idx_correlation ON entity_event (correlation_id) USING BITMAP;

#4.2 事件类型分类

Code
事件类型体系:

系统事件(自动产生):
├── EntityCreated        实体创建
├── EntityUpdated        实体更新(含 diff)
├── EntityDeleted        实体删除
├── EdgeCreated          关系创建
├── EdgeDeleted          关系删除
├── PropertyChanged      属性变更
└── ComputedPropUpdated  计算属性刷新

业务事件(外部输入):
├── Transaction          交易事件
├── Alert                告警事件
├── StatusChange         状态变更
├── Measurement          测量数据
├── Interaction          交互事件
└── CustomEvent          自定义事件

审计事件(操作记录):
├── DataAccess           数据访问
├── PermissionChange     权限变更
├── ExportRequested      数据导出
└── QueryExecuted        查询执行

#4.3 事件写入与查询

Python
# python-sdk/ontology_sdk/events/event_writer.py
from datetime import datetime, timezone
from uuid import uuid4

from pydantic import BaseModel, Field


class OntologyEvent(BaseModel):
    """Ontology 事件模型"""
    event_id: str = Field(default_factory=lambda: str(uuid4()))
    world_id: str
    entity_id: str
    entity_type: str
    event_type: str
    event_subtype: str = ""
    event_source: str
    payload: dict
    summary: str = ""
    severity: str = "info"
    event_time: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
    ingestion_time: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
    caused_by: str = ""
    correlation_id: str = ""


class EventWriter:
    """事件写入器"""

    def __init__(self, doris_connection):
        self._conn = doris_connection

    def write_entity_change_event(
        self,
        world_id: str,
        entity_id: str,
        entity_type: str,
        change_type: str,  # Created, Updated, Deleted
        old_values: dict | None = None,
        new_values: dict | None = None,
        caused_by: str = "system",
    ) -> OntologyEvent:
        """写入实体变更事件"""
        diff = {}
        if old_values and new_values:
            for key in set(list(old_values.keys()) + list(new_values.keys())):
                old_val = old_values.get(key)
                new_val = new_values.get(key)
                if old_val != new_val:
                    diff[key] = {"old": old_val, "new": new_val}

        event = OntologyEvent(
            world_id=world_id,
            entity_id=entity_id,
            entity_type=entity_type,
            event_type=f"Entity{change_type}",
            event_source="ontology-engine",
            payload={
                "change_type": change_type,
                "diff": diff,
                "old_values": old_values or {},
                "new_values": new_values or {},
            },
            summary=f"{entity_type} {entity_id} was {change_type.lower()}d",
            caused_by=caused_by,
        )

        self._insert_event(event)
        return event

    def write_business_event(
        self,
        world_id: str,
        entity_id: str,
        entity_type: str,
        event_type: str,
        payload: dict,
        summary: str = "",
        severity: str = "info",
        event_time: datetime | None = None,
        correlation_id: str = "",
    ) -> OntologyEvent:
        """写入业务事件"""
        event = OntologyEvent(
            world_id=world_id,
            entity_id=entity_id,
            entity_type=entity_type,
            event_type=event_type,
            event_source="business",
            payload=payload,
            summary=summary,
            severity=severity,
            event_time=event_time or datetime.now(timezone.utc),
            correlation_id=correlation_id,
        )

        self._insert_event(event)
        return event

    def _insert_event(self, event: OntologyEvent) -> None:
        """插入事件到 Doris"""
        sql = """
            INSERT INTO entity_event (
                event_id, world_id, entity_id, entity_type,
                event_type, event_subtype, event_source,
                payload, summary, severity,
                event_time, ingestion_time,
                caused_by, correlation_id
            ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
        """
        self._conn.execute(sql, [
            event.event_id, event.world_id, event.entity_id,
            event.entity_type, event.event_type, event.event_subtype,
            event.event_source, event.payload, event.summary,
            event.severity, event.event_time, event.ingestion_time,
            event.caused_by, event.correlation_id,
        ])

#5. 三表之间的查询模式

#5.1 常见查询模式分类

Code
查询模式矩阵:

┌────────────────────┬────────────────────┬──────────────────┐
│    单表查询         │    双表联合         │   三表联合        │
├────────────────────┼────────────────────┼──────────────────┤
│ 按类型查实体        │ 实体 + 关系遍历    │ 实体 + 关系 +     │
│ 全文搜索实体        │ 实体 + 事件时间线  │ 事件时间线         │
│ 按属性过滤          │ 关系 + 事件关联    │ (完整上下文查询) │
│ 语义向量搜索        │                    │                    │
│ 事件时间范围查询    │                    │                    │
│ 事件聚合统计        │                    │                    │
└────────────────────┴────────────────────┴──────────────────┘

#5.2 典型查询示例

SQL
-- 模式 1:实体 360 度视图(三表联合)
-- "给我看某个实体的所有信息:属性、关系、最近事件"

-- 基本信息
SELECT * FROM entity_common
WHERE entity_id = 'person-001' AND world_id = 'world-main';

-- 所有关系(出边和入边)
SELECT
    ee.edge_type,
    ee.direction,
    CASE
        WHEN ee.source_id = 'person-001' THEN ee.target_id
        ELSE ee.source_id
    END AS related_entity_id,
    CASE
        WHEN ee.source_id = 'person-001' THEN ee.target_type
        ELSE ee.source_type
    END AS related_entity_type,
    ee.properties AS edge_properties,
    ee.weight
FROM entity_edge ee
WHERE (ee.source_id = 'person-001' OR ee.target_id = 'person-001')
    AND ee.world_id = 'world-main'
    AND ee.is_deleted = FALSE
ORDER BY ee.created_at DESC;

-- 最近 50 条事件
SELECT
    event_type, summary, severity,
    event_time, payload
FROM entity_event
WHERE entity_id = 'person-001'
    AND world_id = 'world-main'
ORDER BY event_time DESC
LIMIT 50;

-- 模式 2:关系图聚合
-- "每个部门有多少人?各部门之间的合作关系强度?"
SELECT
    JSON_EXTRACT(src.properties, '$.department') AS src_dept,
    JSON_EXTRACT(tgt.properties, '$.department') AS tgt_dept,
    COUNT(*) AS collaboration_count,
    AVG(ee.weight) AS avg_strength
FROM entity_edge ee
JOIN entity_common src ON src.entity_id = ee.source_id
    AND src.world_id = ee.world_id
JOIN entity_common tgt ON tgt.entity_id = ee.target_id
    AND tgt.world_id = ee.world_id
WHERE ee.edge_type = 'Collaborates'
    AND ee.world_id = 'world-main'
GROUP BY src_dept, tgt_dept
ORDER BY collaboration_count DESC;

-- 模式 3:事件驱动的实体发现
-- "最近 24 小时产生告警最多的设备"
SELECT
    ec.entity_id,
    ec.display_name,
    JSON_EXTRACT(ec.properties, '$.device_type') AS device_type,
    COUNT(ev.event_id) AS alert_count,
    MAX(ev.severity) AS max_severity
FROM entity_event ev
JOIN entity_common ec ON ec.entity_id = ev.entity_id
    AND ec.world_id = ev.world_id
WHERE ev.event_type = 'Alert'
    AND ev.world_id = 'world-main'
    AND ev.event_time >= DATE_SUB(NOW(), INTERVAL 24 HOUR)
GROUP BY ec.entity_id, ec.display_name, device_type
ORDER BY alert_count DESC
LIMIT 20;

#6. 在 Doris 上的性能优化

#6.1 分区与分桶策略

Code
分区分桶策略:

entity_common:
  分区:无(数据量适中,全量查询多)
  分桶:HASH(entity_id) 32 桶
  原因:点查和范围查询平衡

entity_edge:
  分区:无
  分桶:HASH(edge_id) 16 桶
  原因:边的数量通常少于实体

entity_event:
  分区:RANGE(event_time) 按天动态分区
  分桶:HASH(entity_id) 32 桶
  原因:
    - 时间范围查询可以裁剪分区
    - 按 entity_id 分桶保证同一实体的事件在同一桶
    - 动态分区自动创建和清理

#6.2 物化视图加速

SQL
-- 物化视图 1:实体类型计数(仪表盘用)
CREATE MATERIALIZED VIEW mv_entity_type_count AS
SELECT
    world_id,
    entity_type,
    COUNT(*) AS entity_count
FROM entity_common
WHERE is_deleted = FALSE
GROUP BY world_id, entity_type;

-- 物化视图 2:每日事件统计
CREATE MATERIALIZED VIEW mv_daily_event_stats AS
SELECT
    world_id,
    entity_type,
    event_type,
    severity,
    DATE(event_time) AS event_date,
    COUNT(*) AS event_count
FROM entity_event
GROUP BY world_id, entity_type, event_type, severity, DATE(event_time);

-- 物化视图 3:关系类型统计
CREATE MATERIALIZED VIEW mv_edge_type_stats AS
SELECT
    world_id,
    edge_type,
    source_type,
    target_type,
    COUNT(*) AS edge_count,
    AVG(weight) AS avg_weight
FROM entity_edge
WHERE is_deleted = FALSE
GROUP BY world_id, edge_type, source_type, target_type;

#6.3 性能基准测试

查询类型数据规模平均延迟P99 延迟QPS
单实体点查100M 实体2 ms8 ms15,000
类型过滤100M 实体45 ms120 ms800
JSON 属性查询100M 实体85 ms250 ms400
一跳关系查询500M 边12 ms35 ms5,000
两跳图遍历500M 边180 ms500 ms200
事件时间范围1B 事件35 ms95 ms1,200
三表联合 360综合120 ms350 ms300
全文搜索100M 实体25 ms75 ms2,000

#7. 在 Iceberg 上的版本化存储

#7.1 双写架构

Code
三表在 Doris + Iceberg 上的双写架构:

┌─────────────┐
│  SDK / API  │
└──────┬──────┘
       │
       ▼
┌──────────────┐
│  Write Path  │
│              │
│  ┌────────┐  │     ┌──────────────────┐
│  │ Doris  │──│────>│ Doris Tables     │  ← 主查询(低延迟)
│  │ Writer │  │     │ entity_common    │
│  └────────┘  │     │ entity_edge      │
│              │     │ entity_event     │
│  ┌────────┐  │     └──────────────────┘
│  │Iceberg │──│────>┌──────────────────┐
│  │ Writer │  │     │ Iceberg Tables   │  ← 版本化(时间旅行)
│  └────────┘  │     │ (via Nessie)     │
│              │     └──────────────────┘
└──────────────┘

#7.2 Schema 演化支持

Python
# 三表模型的 Schema 演化不需要 DDL 变更
# 因为属性存在 JSON 列中

# 旧版本实体(type_version=1)
old_entity = {
    "properties": {
        "name": "张三",
        "email": "zhangsan@example.com"
    }
}

# 新版本实体(type_version=2,增加了 phone 字段)
new_entity = {
    "properties": {
        "name": "张三",
        "email": "zhangsan@example.com",
        "phone": "+86-13800138000"  # 新增字段
    }
}

# SDK 层面通过 type_version 做兼容处理
class EntityMigrator:
    """实体版本迁移器"""

    def migrate_v1_to_v2(self, properties: dict) -> dict:
        """v1 -> v2: 添加 phone 默认值"""
        if "phone" not in properties:
            properties["phone"] = ""
        return properties

#8. 与 SDK 的集成

#8.1 ORM 式访问

Python
# python-sdk/ontology_sdk/orm/entity.py
from ontology_sdk.models import Entity, Edge, Event


# 查询实体
person = Entity.get("person-001", world_id="world-main")
print(person.display_name)
print(person.properties["email"])

# 遍历关系
companies = person.edges("WorksAt").targets()
for company in companies:
    print(f"{person.display_name} works at {company.display_name}")

# 查看时间线
events = person.events(
    event_type="PropertyChanged",
    since="2024-01-01",
    limit=10
)
for event in events:
    print(f"[{event.event_time}] {event.summary}")

# 创建关系
Edge.create(
    source=person,
    target=companies[0],
    edge_type="Manages",
    properties={"team_size": 5}
)

# 写入事件
Event.create(
    entity=person,
    event_type="Promotion",
    payload={"new_level": "Staff Engineer"},
    summary="Promoted to Staff Engineer"
)

#Key Takeaways

  1. 三表模型是 Schema-on-Read 的最佳实践:通过 JSON 属性列实现完全的 Schema 灵活性,新增 Ontology 类型不需要任何 DDL 变更,同时通过索引保证查询性能。

  2. 关系作为一等公民带来图查询能力:entity_edge 表的独立设计支持带权重、带时间、带属性的复杂关系建模,在关系型数据库上实现了图遍历查询。

  3. 事件表实现天然的事件溯源:entity_event 与实体分离存储,支持完整的变更历史追踪、时间线展示和审计合规需求。

  4. 动态分区策略优化时序查询:事件表按天动态分区,自动创建未来分区、清理过期分区,配合 HASH 分桶确保同一实体的事件局部性。

  5. 双写架构兼顾性能和版本化:Doris 提供低延迟查询,Iceberg 通过 Nessie 提供数据版本控制,两者结合满足在线服务和数据治理的双重需求。

#Next Article

下一篇 S3-06《OQL:我们设计的 Ontology 查询语言(语法篇)》 将介绍我们为三表模型量身设计的查询语言——OQL(Ontology Query Language),包括 BNF 语法定义、核心语句、指标展开和图遍历的语法设计。

Tags: #ThreeTableModel #EntityCommon #EntityEdge #EntityEvent #Ontology #SchemaDesign #QueryPatterns #Doris #Iceberg #智策平台 #coomia-dip #数据基座