entity_common / entity_edge / entity_event:三表模型设计
Tags: #ThreeTableModel #Ontology #SchemaDesign #QueryPatterns #EntityModel #智策平台
“系列: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 本体论中的数据可以归纳为三种基本类型:
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 三表模型的设计哲学
三表模型设计原则:
原则 1:类型无关的表结构
所有 Entity 类型共享同一张表,类型区分靠 entity_type 字段
→ 新增 Ontology 类型不需要 DDL 变更
原则 2:属性 JSON 化
每个实体的特有属性存储在 JSON 列中
→ 灵活性 > 查询性能(通过索引弥补)
原则 3:关系显式建模
关系不是"外键",而是一等公民,拥有自己的属性
→ 支持带权重、带时间的关系
原则 4:事件流与状态分离
实体表存"当前状态",事件表存"变化历史"
→ 天然支持事件溯源和时间旅行
#2. entity_common 表:实体的统一载体
#2.1 表结构设计
-- 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 字段设计决策解析
为什么每个字段这样设计:
┌──────────────┬─────────────────────────────────────────────┐
│ 字段 │ 设计决策 │
├──────────────┼─────────────────────────────────────────────┤
│ 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 的内部结构约定
{
"__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 表结构设计
-- 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 关系类型与方向
关系的方向性:
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 图遍历查询模式
-- 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 表结构设计
-- 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 事件类型分类
事件类型体系:
系统事件(自动产生):
├── EntityCreated 实体创建
├── EntityUpdated 实体更新(含 diff)
├── EntityDeleted 实体删除
├── EdgeCreated 关系创建
├── EdgeDeleted 关系删除
├── PropertyChanged 属性变更
└── ComputedPropUpdated 计算属性刷新
业务事件(外部输入):
├── Transaction 交易事件
├── Alert 告警事件
├── StatusChange 状态变更
├── Measurement 测量数据
├── Interaction 交互事件
└── CustomEvent 自定义事件
审计事件(操作记录):
├── DataAccess 数据访问
├── PermissionChange 权限变更
├── ExportRequested 数据导出
└── QueryExecuted 查询执行
#4.3 事件写入与查询
# 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 常见查询模式分类
查询模式矩阵:
┌────────────────────┬────────────────────┬──────────────────┐
│ 单表查询 │ 双表联合 │ 三表联合 │
├────────────────────┼────────────────────┼──────────────────┤
│ 按类型查实体 │ 实体 + 关系遍历 │ 实体 + 关系 + │
│ 全文搜索实体 │ 实体 + 事件时间线 │ 事件时间线 │
│ 按属性过滤 │ 关系 + 事件关联 │ (完整上下文查询) │
│ 语义向量搜索 │ │ │
│ 事件时间范围查询 │ │ │
│ 事件聚合统计 │ │ │
└────────────────────┴────────────────────┴──────────────────┘
#5.2 典型查询示例
-- 模式 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 分区与分桶策略
分区分桶策略:
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 物化视图加速
-- 物化视图 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 ms | 8 ms | 15,000 |
| 类型过滤 | 100M 实体 | 45 ms | 120 ms | 800 |
| JSON 属性查询 | 100M 实体 | 85 ms | 250 ms | 400 |
| 一跳关系查询 | 500M 边 | 12 ms | 35 ms | 5,000 |
| 两跳图遍历 | 500M 边 | 180 ms | 500 ms | 200 |
| 事件时间范围 | 1B 事件 | 35 ms | 95 ms | 1,200 |
| 三表联合 360 | 综合 | 120 ms | 350 ms | 300 |
| 全文搜索 | 100M 实体 | 25 ms | 75 ms | 2,000 |
#7. 在 Iceberg 上的版本化存储
#7.1 双写架构
三表在 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 演化支持
# 三表模型的 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-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
-
三表模型是 Schema-on-Read 的最佳实践:通过 JSON 属性列实现完全的 Schema 灵活性,新增 Ontology 类型不需要任何 DDL 变更,同时通过索引保证查询性能。
-
关系作为一等公民带来图查询能力:entity_edge 表的独立设计支持带权重、带时间、带属性的复杂关系建模,在关系型数据库上实现了图遍历查询。
-
事件表实现天然的事件溯源:entity_event 与实体分离存储,支持完整的变更历史追踪、时间线展示和审计合规需求。
-
动态分区策略优化时序查询:事件表按天动态分区,自动创建未来分区、清理过期分区,配合 HASH 分桶确保同一实体的事件局部性。
-
双写架构兼顾性能和版本化: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 #数据基座