返回博客

分析引擎:OLAP 能力的 Ontology 封装

Tags: #AnalyticsEngine #OLAP #Doris #Aggregation #Dashboard #智策平台

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

系列: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

Code
传统 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 架构定位

Code
分析引擎在数据基座中的位置:

┌──────────────────────────────────────────┐
│              OQL Query Interface          │
├──────────────────────────────────────────┤
│  ┌────────┐  ┌────────┐  ┌────────────┐ │
│  │ FETCH  │  │TRAVERSE│  │ AGGREGATE  │ │
│  │ 实体查询│  │ 图遍历  │  │ 分析引擎   │ │
│  └────────┘  └────────┘  └────────────┘ │
├──────────────────────────────────────────┤
│         Query Optimization Layer          │
│  (谓词下推、物化视图匹配、聚合策略选择)      │
├──────────────────────────────────────────┤
│         Apache Doris OLAP Engine          │
│  (列式存储、向量化执行、物化视图)            │
└──────────────────────────────────────────┘

#2. AGGREGATE 语句编译

#2.1 OQL AGGREGATE 语法

BNF
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

Python
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

Python
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 多维分析示例

Code
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 指标定义

Python
@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 指标展开机制

Python
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 自动创建策略

Python
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 物化视图特性

Code
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 数据推送架构

Code
实时仪表盘数据推送:

┌──────────┐   ┌──────────┐   ┌──────────┐
│ Flink CDC │──→│ Doris    │──→│ MV 刷新   │
│ (实时入库) │   │ (存储)    │   │ (自动)    │
└──────────┘   └──────────┘   └──────────┘
                                    │
                                    ▼
                              ┌──────────┐
                              │ Dashboard │
                              │ Query     │
                              │ Service   │
                              └──────────┘
                                    │
                    ┌───────────────┼───────────────┐
                    │               │               │
                    ▼               ▼               ▼
              ┌──────────┐  ┌──────────┐  ┌──────────┐
              │ WebSocket│  │ SSE      │  │ Polling  │
              │ 推送      │  │ 推送     │  │ 轮询     │
              └──────────┘  └──────────┘  └──────────┘

#6.2 仪表盘查询优化

Python
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. 性能基准

Code
分析引擎性能基准:

测试环境: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. 测试策略

Python
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

  1. Ontology 语义的分析查询大幅降低了 BI 门槛:用户通过 OQL AGGREGATE 以实体类型和属性为维度分析,无需理解底层表结构和 SQL。

  2. 物化视图是分析性能的决定性因素:自动创建的物化视图可以将聚合查询从秒级降到毫秒级(50-160 倍提升)。

  3. 指标系统实现了指标定义与查询的解耦:一次定义、处处展开——指标公式在 Metric Registry 中注册,通过 WITH METRIC 子句在任意查询中引用。

  4. 多维分析(CUBE/ROLLUP)覆盖了所有分析维度组合:通过 Doris 原生 CUBE/ROLLUP 支持,用户一条查询即可获得所有维度组合的聚合结果。

  5. 实时仪表盘依赖分层缓存和物化视图:MV 提供毫秒级基础查询,配合结果缓存和 WebSocket 推送实现实时数据展示。

#Next Article

下一篇 S3-13《搜索引擎:Elasticsearch 的 Ontology 集成》 将展示如何将 Elasticsearch 的全文搜索能力集成到 Ontology 查询中。

Tags: #AnalyticsEngine #OLAP #Doris #Aggregation #MaterializedView #MetricSystem #CUBE #ROLLUP #Dashboard #智策平台 #coomia-dip #数据基座