分析引擎:OLAP 能力的 Ontology 封装
Tags: #AnalyticsEngine #OLAP #Doris #Aggregation #Dashboard #智策平台
“系列:S3 数据基座 · 第 12 篇 | 难度:高级 | 阅读时间:20 分钟
分析引擎:OLAP 能力的 Ontology 封装
Tags: #AnalyticsEngine #OLAP #Doris #Aggregation #Dashboard #智策平台
#TL;DR
coomia-dip 平台将 Apache Doris 的 OLAP 分析能力封装为 Ontology 语义的分析查询接口。用户不再编写复杂的 SQL 聚合查询,而是通过 OQL 的 AGGREGATE 语句以实体类型和属性为维度进行分析。本文完整解析分析引擎的架构设计,包括 OQL AGGREGATE 到 Doris SQL 的编译流程、多维度聚合(Cube / Rollup)实现、预聚合物化视图自动创建与匹配、指标系统(Metric System)的定义与展开机制、实时仪表盘数据推送和分析查询的性能优化策略。
#1. 分析引擎定位
#1.1 Ontology 分析 vs 传统 BI
传统 BI 分析 vs Ontology 分析:
传统 BI:
用户 → SQL 查询 → 数据仓库 → 结果
痛点:用户需要理解表结构、写 SQL
Ontology 分析:
用户 → OQL AGGREGATE → 分析引擎 → 结果
优势:用户只需理解实体类型和属性
示例对比:
SQL: SELECT department, AVG(age), COUNT(*)
FROM entity_common
WHERE entity_type = 'Person'
AND JSON_EXTRACT(properties, '$.status') = 'active'
GROUP BY JSON_EXTRACT(properties, '$.department')
OQL: AGGREGATE Person
GROUP BY department
SELECT AVG(age), COUNT(*)
WHERE status = 'active'
#1.2 架构定位
分析引擎在数据基座中的位置:
┌──────────────────────────────────────────┐
│ OQL Query Interface │
├──────────────────────────────────────────┤
│ ┌────────┐ ┌────────┐ ┌────────────┐ │
│ │ FETCH │ │TRAVERSE│ │ AGGREGATE │ │
│ │ 实体查询│ │ 图遍历 │ │ 分析引擎 │ │
│ └────────┘ └────────┘ └────────────┘ │
├──────────────────────────────────────────┤
│ Query Optimization Layer │
│ (谓词下推、物化视图匹配、聚合策略选择) │
├──────────────────────────────────────────┤
│ Apache Doris OLAP Engine │
│ (列式存储、向量化执行、物化视图) │
└──────────────────────────────────────────┘
#2. AGGREGATE 语句编译
#2.1 OQL AGGREGATE 语法
aggregate_stmt ::= 'AGGREGATE' entity_type
'GROUP' 'BY' group_columns
'SELECT' agg_expressions
['WHERE' predicate]
['HAVING' having_predicate]
['ORDER' 'BY' order_columns]
['LIMIT' number]
agg_expressions ::= agg_expr (',' agg_expr)*
agg_expr ::= agg_function '(' column ')' ['AS' alias]
agg_function ::= 'COUNT' | 'SUM' | 'AVG' | 'MIN' | 'MAX'
| 'PERCENTILE' | 'STDDEV' | 'VARIANCE'
| 'DISTINCT_COUNT' | 'TOP_N'
#2.2 编译为 Doris SQL
class AggregateCompiler:
"""AGGREGATE 语句到 Doris SQL 的编译器"""
def compile(self, stmt: AggregateStatement) -> str:
entity_type = stmt.entity_type
schema = self._registry.get_entity_type(entity_type)
# 构建 SELECT 子句
select_parts = []
for group_col in stmt.group_by:
prop = schema.properties[group_col]
select_parts.append(
f"JSON_EXTRACT(properties, '$.{group_col}') AS {group_col}"
)
for agg in stmt.aggregations:
sql_agg = self._compile_aggregation(agg, schema)
select_parts.append(sql_agg)
# 构建 WHERE 子句
where_parts = [f"entity_type = '{entity_type}'"]
if stmt.where_clause:
where_parts.append(
self._compile_predicate(stmt.where_clause, schema)
)
# 构建 GROUP BY 子句
group_parts = [
f"JSON_EXTRACT(properties, '$.{col}')"
for col in stmt.group_by
]
sql = f"""
SELECT {', '.join(select_parts)}
FROM entity_common
WHERE {' AND '.join(where_parts)}
GROUP BY {', '.join(group_parts)}
"""
if stmt.having_clause:
sql += f" HAVING {self._compile_predicate(stmt.having_clause, schema)}"
if stmt.order_by:
order_parts = [
f"{col} {direction}" for col, direction in stmt.order_by
]
sql += f" ORDER BY {', '.join(order_parts)}"
if stmt.limit:
sql += f" LIMIT {stmt.limit}"
return sql
def _compile_aggregation(
self, agg: AggregationExpr, schema: EntitySchema
) -> str:
col = agg.column
func = agg.function.upper()
alias = agg.alias or f"{func.lower()}_{col}"
col_expr = f"JSON_EXTRACT(properties, '$.{col}')"
if func == 'DISTINCT_COUNT':
return f"COUNT(DISTINCT {col_expr}) AS {alias}"
elif func == 'TOP_N':
return f"TOPN({col_expr}, {agg.n}) AS {alias}"
elif func == 'PERCENTILE':
return f"PERCENTILE_APPROX({col_expr}, {agg.percentile}) AS {alias}"
else:
return f"{func}({col_expr}) AS {alias}"
#3. 多维分析
#3.1 Cube 和 Rollup
class MultiDimensionalAnalysis:
"""多维分析引擎"""
def cube_query(
self,
entity_type: str,
dimensions: list[str],
measures: list[AggregationExpr]
) -> str:
"""生成 CUBE 查询 — 所有维度组合的聚合"""
schema = self._registry.get_entity_type(entity_type)
dim_exprs = [
f"JSON_EXTRACT(properties, '$.{d}') AS {d}"
for d in dimensions
]
measure_exprs = [
self._compile_aggregation(m, schema) for m in measures
]
return f"""
SELECT {', '.join(dim_exprs + measure_exprs)}
FROM entity_common
WHERE entity_type = '{entity_type}'
GROUP BY CUBE({', '.join(
f"JSON_EXTRACT(properties, '$.{d}')" for d in dimensions
)})
"""
def rollup_query(
self,
entity_type: str,
hierarchy: list[str],
measures: list[AggregationExpr]
) -> str:
"""生成 ROLLUP 查询 — 层次聚合"""
schema = self._registry.get_entity_type(entity_type)
dim_exprs = [
f"JSON_EXTRACT(properties, '$.{d}') AS {d}"
for d in hierarchy
]
measure_exprs = [
self._compile_aggregation(m, schema) for m in measures
]
return f"""
SELECT {', '.join(dim_exprs + measure_exprs)}
FROM entity_common
WHERE entity_type = '{entity_type}'
GROUP BY ROLLUP({', '.join(
f"JSON_EXTRACT(properties, '$.{d}')" for d in hierarchy
)})
"""
#3.2 多维分析示例
CUBE 分析示例:
OQL:
AGGREGATE Device
CUBE BY region, type, status
SELECT COUNT(*), AVG(uptime)
结果(8 种维度组合):
┌──────────┬──────────┬──────────┬───────┬──────────┐
│ region │ type │ status │ count │ avg_up │
├──────────┼──────────┼──────────┼───────┼──────────┤
│ 华东 │ 传感器 │ 运行中 │ 150 │ 99.2% │
│ 华东 │ 传感器 │ NULL │ 200 │ 95.1% │
│ 华东 │ NULL │ 运行中 │ 300 │ 98.5% │
│ 华东 │ NULL │ NULL │ 500 │ 94.2% │
│ NULL │ 传感器 │ 运行中 │ 400 │ 99.0% │
│ NULL │ 传感器 │ NULL │ 600 │ 96.3% │
│ NULL │ NULL │ 运行中 │ 800 │ 98.8% │
│ NULL │ NULL │ NULL │ 1200 │ 95.5% │
└──────────┴──────────┴──────────┴───────┴──────────┘
(NULL 表示该维度被聚合)
#4. 指标系统(Metric System)
#4.1 指标定义
@dataclass
class MetricDefinition:
"""指标定义"""
name: str
display_name: str
description: str
entity_type: str
formula: MetricFormula
dimensions: list[str]
filters: list[MetricFilter] | None
time_granularity: str | None # 'DAILY', 'HOURLY', 'MONTHLY'
cache_ttl: timedelta | None
def compile_to_sql(self, schema: EntitySchema) -> str:
"""编译为 SQL 子查询"""
return self.formula.compile(schema)
@dataclass
class MetricFormula:
"""指标计算公式"""
type: str # 'SIMPLE', 'DERIVED', 'COMPOSITE'
expression: str
def compile(self, schema: EntitySchema) -> str:
if self.type == 'SIMPLE':
return self._compile_simple(schema)
elif self.type == 'DERIVED':
return self._compile_derived(schema)
elif self.type == 'COMPOSITE':
return self._compile_composite(schema)
# 指标注册示例
metric_registry.register(MetricDefinition(
name='equipment_utilization',
display_name='设备利用率',
description='设备运行时间 / 总时间',
entity_type='Device',
formula=MetricFormula(
type='DERIVED',
expression='SUM(running_hours) / SUM(total_hours) * 100'
),
dimensions=['region', 'device_type'],
time_granularity='DAILY',
cache_ttl=timedelta(minutes=5),
))
metric_registry.register(MetricDefinition(
name='direct_reports',
display_name='直接下属数',
description='通过 Manages 关系连接的下属实体数量',
entity_type='Person',
formula=MetricFormula(
type='SIMPLE',
expression='COUNT(edge WHERE edge_type = "Manages")'
),
dimensions=['department'],
cache_ttl=timedelta(minutes=30),
))
#4.2 指标展开机制
class MetricExpansionEngine:
"""指标展开引擎"""
def expand_metric(
self, entity_type: str, metric_name: str, context: QueryContext
) -> str:
"""将指标名展开为 SQL 子查询"""
metric = self._registry.get_metric(metric_name)
if metric is None:
raise MetricNotFoundError(f"Unknown metric: {metric_name}")
if metric.entity_type != entity_type:
raise MetricTypeMismatchError(
f"Metric '{metric_name}' is defined on "
f"'{metric.entity_type}', not '{entity_type}'"
)
# 检查缓存
if metric.cache_ttl:
cached = self._cache.get(metric_name, context)
if cached:
return cached
# 编译指标公式
if metric.formula.type == 'SIMPLE':
sql = self._expand_simple(metric, context)
elif metric.formula.type == 'DERIVED':
sql = self._expand_derived(metric, context)
elif metric.formula.type == 'COMPOSITE':
sql = self._expand_composite(metric, context)
return sql
def _expand_simple(
self, metric: MetricDefinition, context: QueryContext
) -> str:
"""展开简单指标(基于边关系的计数/聚合)"""
if 'edge' in metric.formula.expression.lower():
# 基于边关系的指标
edge_type = self._extract_edge_type(metric.formula.expression)
agg_func = self._extract_agg_func(metric.formula.expression)
return f"""
(SELECT {agg_func}(*)
FROM entity_edge
WHERE source_id = t.entity_id
AND edge_type = '{edge_type}')
"""
return metric.formula.compile(context.schema)
#5. 预聚合物化视图
#5.1 自动创建策略
class AutoMaterializedViewManager:
"""自动物化视图管理器"""
def analyze_and_create(self, query_log: list[QueryLogEntry]):
"""分析查询日志,自动创建有价值的物化视图"""
# 提取聚合模式
patterns = self._extract_aggregation_patterns(query_log)
for pattern in patterns:
if self._should_materialize(pattern):
self._create_materialized_view(pattern)
def _should_materialize(self, pattern: AggregationPattern) -> bool:
"""判断是否值得创建物化视图"""
# 条件 1:每天执行 > 100 次
if pattern.daily_executions < 100:
return False
# 条件 2:平均执行时间 > 500ms
if pattern.avg_execution_ms < 500:
return False
# 条件 3:结果集 < 100MB
if pattern.estimated_result_size_mb > 100:
return False
# 条件 4:基础表更新频率不高于每分钟
if pattern.source_update_frequency_per_minute > 1:
return False
return True
def _create_materialized_view(self, pattern: AggregationPattern):
"""创建物化视图"""
mv_name = f"mv_{pattern.entity_type}_{pattern.hash[:8]}"
dim_cols = [
f"JSON_EXTRACT(properties, '$.{d}') AS {d}"
for d in pattern.dimensions
]
agg_cols = [
f"{a.function}(JSON_EXTRACT(properties, '$.{a.column}')) AS {a.alias}"
for a in pattern.aggregations
]
sql = f"""
CREATE MATERIALIZED VIEW {mv_name} AS
SELECT
entity_type,
{', '.join(dim_cols)},
{', '.join(agg_cols)}
FROM entity_common
WHERE entity_type = '{pattern.entity_type}'
GROUP BY entity_type, {', '.join(
f"JSON_EXTRACT(properties, '$.{d}')" for d in pattern.dimensions
)}
"""
self._doris.execute(sql)
#5.2 Doris 物化视图特性
Doris 物化视图自动匹配:
原始查询:
SELECT region, COUNT(*), AVG(uptime)
FROM entity_common
WHERE entity_type = 'Device'
GROUP BY JSON_EXTRACT(properties, '$.region')
物化视图定义:
CREATE MATERIALIZED VIEW mv_device_region AS
SELECT entity_type,
JSON_EXTRACT(properties, '$.region') AS region,
COUNT(*) AS cnt,
SUM(JSON_EXTRACT(properties, '$.uptime')) AS sum_uptime,
COUNT(JSON_EXTRACT(properties, '$.uptime')) AS cnt_uptime
FROM entity_common
GROUP BY entity_type, JSON_EXTRACT(properties, '$.region')
Doris 自动匹配过程:
1. 检测到 GROUP BY 维度匹配
2. AVG(uptime) 可从 SUM(uptime)/COUNT(uptime) 推导
3. 自动改写查询使用物化视图
改写后:
SELECT region, cnt, sum_uptime / cnt_uptime
FROM mv_device_region
WHERE entity_type = 'Device'
效果:查询时间从 2s → 20ms(100x 提升)
#6. 实时仪表盘
#6.1 数据推送架构
实时仪表盘数据推送:
┌──────────┐ ┌──────────┐ ┌──────────┐
│ Flink CDC │──→│ Doris │──→│ MV 刷新 │
│ (实时入库) │ │ (存储) │ │ (自动) │
└──────────┘ └──────────┘ └──────────┘
│
▼
┌──────────┐
│ Dashboard │
│ Query │
│ Service │
└──────────┘
│
┌───────────────┼───────────────┐
│ │ │
▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌──────────┐
│ WebSocket│ │ SSE │ │ Polling │
│ 推送 │ │ 推送 │ │ 轮询 │
└──────────┘ └──────────┘ └──────────┘
#6.2 仪表盘查询优化
class DashboardQueryService:
"""仪表盘查询服务"""
async def get_dashboard_data(
self, dashboard_id: str
) -> DashboardData:
dashboard = self._registry.get_dashboard(dashboard_id)
panels = dashboard.panels
# 并行执行所有面板查询
tasks = [
self._execute_panel_query(panel) for panel in panels
]
results = await asyncio.gather(*tasks)
return DashboardData(
dashboard_id=dashboard_id,
panels={
panel.id: result
for panel, result in zip(panels, results)
},
timestamp=datetime.utcnow()
)
async def _execute_panel_query(
self, panel: DashboardPanel
) -> PanelData:
# 优先使用物化视图
mv = self._mv_matcher.find_matching_mv(panel.query)
if mv:
return await self._execute_from_mv(mv, panel)
# 退而使用缓存
cache_key = panel.query_hash
cached = self._cache.get(cache_key)
if cached and cached.age < panel.refresh_interval:
return cached
# 执行实时查询
result = await self._execute_query(panel.query)
self._cache.set(cache_key, result, ttl=panel.refresh_interval)
return result
#7. 性能基准
分析引擎性能基准:
测试环境:100 万实体,5 种类型,3 Doris BE 节点
┌─────────────────────┬──────────┬──────────┬──────────┐
│ 查询类型 │ 无优化 │ 有 MV │ 提升 │
├─────────────────────┼──────────┼──────────┼──────────┤
│ 单维度 COUNT │ 800ms │ 15ms │ 53x │
│ 单维度 AVG │ 1.2s │ 20ms │ 60x │
│ 两维度 GROUP BY │ 2.5s │ 30ms │ 83x │
│ CUBE (3维度) │ 8.0s │ 50ms │ 160x │
│ ROLLUP (4层级) │ 5.0s │ 40ms │ 125x │
│ PERCENTILE_APPROX │ 3.0s │ 100ms │ 30x │
│ 跨类型聚合 │ 4.0s │ 200ms │ 20x │
└─────────────────────┴──────────┴──────────┴──────────┘
结论:物化视图是分析查询性能的决定性因素
#8. 测试策略
class TestAnalyticsEngine:
def test_aggregate_compilation(self):
oql = """
AGGREGATE Person
GROUP BY department
SELECT COUNT(*), AVG(age)
WHERE status = 'active'
"""
sql = compiler.compile(parse(oql))
assert "GROUP BY" in sql
assert "COUNT(*)" in sql
assert "AVG" in sql
def test_metric_expansion(self):
result = execute(
"FETCH Person WITH METRIC direct_reports LIMIT 10"
)
assert 'direct_reports' in result.columns
assert all(isinstance(v, int) for v in result['direct_reports'])
def test_mv_auto_matching(self):
# 创建物化视图
create_mv("mv_test", "SELECT type, COUNT(*) FROM entity_common GROUP BY type")
# 查询应自动匹配
plan = explain("AGGREGATE Device GROUP BY type SELECT COUNT(*)")
assert "mv_test" in plan.used_views
def test_cube_correctness(self):
result = execute("""
AGGREGATE Device CUBE BY region, type SELECT COUNT(*)
""")
# CUBE 应产生 2^2 = 4 种组合
null_rows = [r for r in result if r['region'] is None or r['type'] is None]
assert len(null_rows) > 0 # 应有聚合行
#Key Takeaways
-
Ontology 语义的分析查询大幅降低了 BI 门槛:用户通过 OQL AGGREGATE 以实体类型和属性为维度分析,无需理解底层表结构和 SQL。
-
物化视图是分析性能的决定性因素:自动创建的物化视图可以将聚合查询从秒级降到毫秒级(50-160 倍提升)。
-
指标系统实现了指标定义与查询的解耦:一次定义、处处展开——指标公式在 Metric Registry 中注册,通过 WITH METRIC 子句在任意查询中引用。
-
多维分析(CUBE/ROLLUP)覆盖了所有分析维度组合:通过 Doris 原生 CUBE/ROLLUP 支持,用户一条查询即可获得所有维度组合的聚合结果。
-
实时仪表盘依赖分层缓存和物化视图:MV 提供毫秒级基础查询,配合结果缓存和 WebSocket 推送实现实时数据展示。
#Next Article
下一篇 S3-13《搜索引擎:Elasticsearch 的 Ontology 集成》 将展示如何将 Elasticsearch 的全文搜索能力集成到 Ontology 查询中。
Tags: #AnalyticsEngine #OLAP #Doris #Aggregation #MaterializedView #MetricSystem #CUBE #ROLLUP #Dashboard #智策平台 #coomia-dip #数据基座