返回博客

S3-15 物化视图自动化:注册指标即创建物化视图

智策平台的 MaterializedViewService 实现了"注册指标即创建物化视图"的自动化能力。支持三种刷新模式(手动 / 定时 / 事件驱动),提供过期检测机制,并在 Schema 变更时自动级联失效。本文完整拆解物化视图的自动创建、刷新调度、过期检测和 Schema 变更级联的全链路实现。

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

S3-15 物化视图自动化:注册指标即创建物化视图

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

#TL;DR

智策平台的 MaterializedViewService 实现了"注册指标即创建物化视图"的自动化能力。支持三种刷新模式(手动 / 定时 / 事件驱动),提供过期检测机制,并在 Schema 变更时自动级联失效。本文完整拆解物化视图的自动创建、刷新调度、过期检测和 Schema 变更级联的全链路实现。

#1. 为什么需要物化视图自动化

物化视图(Materialized View, MV)是预计算查询结果并持久化存储的技术。在传统数据库中,创建和维护 MV 需要 DBA 手工操作:编写 CREATE MATERIALIZED VIEW 语句、配置刷新策略、监控数据新鲜度、处理表结构变更后的 MV 重建。

在 Ontology 平台中,指标(Metric)是动态注册的。运营人员可能每天新增或修改指标定义。如果每个指标的物化视图都需要 DBA 手工维护,这个流程根本无法扩展。

智策平台的设计目标是:注册指标时自动创建物化视图,删除指标时自动销毁,Schema 变更时自动重建。

Code
+------------------------------------------------------------------+
|  Metric Registration Flow with Auto-MV                            |
|                                                                   |
|  MetricRegistryService.register(metric_def)                       |
|       |                                                           |
|       v                                                           |
|  strategies includes MATERIALIZED?                                |
|       |YES                                                        |
|       v                                                           |
|  MaterializedViewService.create_for_metric(metric_def)            |
|       |                                                           |
|       +---> Generate CREATE MV SQL                                |
|       +---> Execute on Doris                                      |
|       +---> Register refresh schedule                             |
|       +---> Register staleness monitor                            |
|       +---> Register schema-change listener                      |
+------------------------------------------------------------------+

#2. MaterializedViewService 架构

Code
+------------------------------------------------------------------+
|                  MaterializedViewService                          |
|                                                                   |
|  +--------------------+  +---------------------+                  |
|  | MV Creator         |  | MV Metadata Store   |                  |
|  | - generate DDL     |  | - PostgreSQL        |                  |
|  | - execute on Doris |  | - MV <-> Metric map |                  |
|  +--------------------+  +---------------------+                  |
|           |                        |                              |
|  +--------------------+  +---------------------+                  |
|  | Refresh Scheduler  |  | Staleness Detector  |                  |
|  | - manual           |  | - last_refresh_at   |                  |
|  | - scheduled (cron) |  | - freshness_threshold|                 |
|  | - event-driven     |  | - alert on stale    |                  |
|  +--------------------+  +---------------------+                  |
|           |                        |                              |
|  +--------------------+  +---------------------+                  |
|  | Schema Change      |  | MV Health Monitor   |                  |
|  | Listener           |  | - row count         |                  |
|  | - invalidate       |  | - refresh duration  |                  |
|  | - rebuild          |  | - error tracking    |                  |
|  +--------------------+  +---------------------+                  |
+------------------------------------------------------------------+

#3. MV 元数据模型

#3.1 MaterializedViewDefinition

Python
class MaterializedViewDefinition(BaseModel):
    """物化视图定义"""
    mv_id: str = Field(default_factory=lambda: str(uuid4()))
    mv_name: str                       # mv_customer_order_count
    metric_id: str                     # 关联的指标 ID
    metric_name: str                   # 关联的指标名称
    object_type: str                   # Customer
    source_tables: list[str]           # 源表列表
    create_sql: str                    # CREATE MV 语句
    query_sql: str                     # 内部查询 SQL
    refresh_mode: RefreshMode          # 刷新模式
    refresh_schedule: str | None = None # cron 表达式
    status: MvStatus = MvStatus.CREATING
    row_count: int = 0
    last_refresh_at: datetime | None = None
    last_refresh_duration_ms: int = 0
    freshness_threshold_seconds: int = 3600
    created_at: datetime = Field(default_factory=datetime.utcnow)
    updated_at: datetime = Field(default_factory=datetime.utcnow)
    error_message: str | None = None

class RefreshMode(str, Enum):
    MANUAL = "MANUAL"           # 手动刷新
    SCHEDULED = "SCHEDULED"     # 定时刷新
    EVENT_DRIVEN = "EVENT_DRIVEN"  # 事件驱动刷新

class MvStatus(str, Enum):
    CREATING = "CREATING"
    ACTIVE = "ACTIVE"
    REFRESHING = "REFRESHING"
    STALE = "STALE"
    ERROR = "ERROR"
    REBUILDING = "REBUILDING"
    DROPPED = "DROPPED"

#4. 自动创建物化视图

#4.1 从指标定义生成 MV

Python
class MvCreator:
    """物化视图创建器"""

    async def create_for_metric(
        self, metric: MetricDefinition
    ) -> MaterializedViewDefinition:
        """根据指标定义自动创建物化视图"""

        # 1. 生成 MV 名称
        mv_name = f"mv_{metric.name}"

        # 2. 生成查询 SQL
        query_sql = self._generate_query_sql(metric)

        # 3. 生成 CREATE MV SQL
        create_sql = self._generate_create_sql(mv_name, query_sql, metric)

        # 4. 确定刷新模式
        refresh_mode = self._determine_refresh_mode(metric)
        refresh_schedule = self._determine_schedule(metric)

        # 5. 在 Doris 上执行创建
        try:
            await self._doris.execute(create_sql)
        except DorisError as e:
            raise MvCreationError(
                f"Failed to create MV {mv_name}: {e}"
            )

        # 6. 保存元数据
        mv_def = MaterializedViewDefinition(
            mv_name=mv_name,
            metric_id=metric.metric_id,
            metric_name=metric.name,
            object_type=metric.object_type,
            source_tables=self._extract_source_tables(query_sql),
            create_sql=create_sql,
            query_sql=query_sql,
            refresh_mode=refresh_mode,
            refresh_schedule=refresh_schedule,
            status=MvStatus.ACTIVE,
            freshness_threshold_seconds=self._freshness_to_seconds(
                metric.freshness_requirement
            ),
        )
        await self._metadata_store.save(mv_def)

        # 7. 注册刷新调度
        if refresh_mode == RefreshMode.SCHEDULED:
            await self._scheduler.register(mv_def)
        elif refresh_mode == RefreshMode.EVENT_DRIVEN:
            await self._event_listener.register(mv_def)

        # 8. 触发初始刷新
        await self._trigger_refresh(mv_def)

        return mv_def

    def _generate_query_sql(self, metric: MetricDefinition) -> str:
        """根据指标表达式生成查询 SQL"""
        expr = metric.expression
        obj_table = f"{metric.object_type.lower()}_objects"

        if expr.type == ExpressionType.SIMPLE_AGG:
            agg = expr.aggregation.value
            source = expr.source_field
            filter_clause = (
                f"WHERE {expr.filter_clause}" if expr.filter_clause else ""
            )
            return f"""
            SELECT
                o.object_id,
                {agg}(r.{source}) AS metric_value,
                NOW() AS computed_at
            FROM {obj_table} o
            LEFT JOIN {source}_objects r
                ON r.{metric.object_type.lower()}_id = o.object_id
            {filter_clause}
            GROUP BY o.object_id
            """

        elif expr.type == ExpressionType.SQL:
            # 将参数化 SQL 转为聚合形式
            return self._wrap_as_mv_query(expr.sql, metric)

        raise MvCreationError(
            f"Cannot create MV for expression type: {expr.type}"
        )

    def _generate_create_sql(
        self, mv_name: str, query_sql: str, metric: MetricDefinition
    ) -> str:
        """生成 Doris CREATE MATERIALIZED VIEW 语句"""
        return f"""
        CREATE MATERIALIZED VIEW IF NOT EXISTS {mv_name}
        BUILD DEFERRED
        REFRESH ASYNC
        DISTRIBUTED BY HASH(object_id) BUCKETS AUTO
        AS
        {query_sql}
        """

    def _determine_refresh_mode(
        self, metric: MetricDefinition
    ) -> RefreshMode:
        """根据新鲜度要求确定刷新模式"""
        freshness = metric.freshness_requirement
        if freshness in (FreshnessLevel.REALTIME, FreshnessLevel.NEAR_REALTIME):
            return RefreshMode.EVENT_DRIVEN
        elif freshness in (FreshnessLevel.HOURLY, FreshnessLevel.DAILY):
            return RefreshMode.SCHEDULED
        else:
            return RefreshMode.MANUAL

    def _determine_schedule(
        self, metric: MetricDefinition
    ) -> str | None:
        """生成 cron 表达式"""
        freshness = metric.freshness_requirement
        schedules = {
            FreshnessLevel.HOURLY: "0 * * * *",      # 每小时
            FreshnessLevel.DAILY: "0 2 * * *",        # 每天凌晨 2 点
            FreshnessLevel.WEEKLY: "0 2 * * 1",       # 每周一凌晨 2 点
        }
        return schedules.get(freshness)

    def _freshness_to_seconds(self, freshness: FreshnessLevel) -> int:
        mapping = {
            FreshnessLevel.REALTIME: 60,
            FreshnessLevel.NEAR_REALTIME: 300,
            FreshnessLevel.HOURLY: 7200,     # 2 小时阈值
            FreshnessLevel.DAILY: 172800,     # 2 天阈值
            FreshnessLevel.WEEKLY: 1209600,   # 2 周阈值
        }
        return mapping.get(freshness, 86400)

#5. 三种刷新模式

#5.1 手动刷新(MANUAL)

手动刷新由用户或管理员显式触发:

Python
class ManualRefreshHandler:
    """手动刷新处理器"""

    async def refresh(self, mv_def: MaterializedViewDefinition) -> RefreshResult:
        mv_def.status = MvStatus.REFRESHING
        await self._metadata_store.save(mv_def)

        start = time.monotonic()
        try:
            await self._doris.execute(
                f"REFRESH MATERIALIZED VIEW {mv_def.mv_name}"
            )

            elapsed_ms = int((time.monotonic() - start) * 1000)
            row_count = await self._get_row_count(mv_def.mv_name)

            mv_def.status = MvStatus.ACTIVE
            mv_def.last_refresh_at = datetime.utcnow()
            mv_def.last_refresh_duration_ms = elapsed_ms
            mv_def.row_count = row_count
            mv_def.error_message = None
            await self._metadata_store.save(mv_def)

            return RefreshResult(
                success=True, duration_ms=elapsed_ms, row_count=row_count
            )

        except Exception as e:
            mv_def.status = MvStatus.ERROR
            mv_def.error_message = str(e)
            await self._metadata_store.save(mv_def)
            return RefreshResult(success=False, error=str(e))

#5.2 定时刷新(SCHEDULED)

Python
class ScheduledRefreshHandler:
    """定时刷新处理器"""

    def __init__(self, scheduler: AsyncScheduler):
        self._scheduler = scheduler
        self._refresh_handler = ManualRefreshHandler()

    async def register(self, mv_def: MaterializedViewDefinition) -> None:
        """注册定时刷新任务"""
        job_id = f"mv_refresh_{mv_def.mv_name}"

        await self._scheduler.add_job(
            job_id=job_id,
            cron=mv_def.refresh_schedule,
            callback=self._on_scheduled_refresh,
            args={"mv_id": mv_def.mv_id},
        )
        logger.info(
            f"Registered refresh schedule for {mv_def.mv_name}: "
            f"{mv_def.refresh_schedule}"
        )

    async def unregister(self, mv_def: MaterializedViewDefinition) -> None:
        job_id = f"mv_refresh_{mv_def.mv_name}"
        await self._scheduler.remove_job(job_id)

    async def _on_scheduled_refresh(self, mv_id: str) -> None:
        """定时触发的刷新回调"""
        mv_def = await self._metadata_store.get(mv_id)
        if mv_def is None or mv_def.status == MvStatus.DROPPED:
            return

        if mv_def.status == MvStatus.REFRESHING:
            logger.warning(
                f"Skipping scheduled refresh for {mv_def.mv_name}: "
                "already refreshing"
            )
            return

        result = await self._refresh_handler.refresh(mv_def)
        if not result.success:
            await self._alert_service.send(
                AlertLevel.WARNING,
                f"MV refresh failed: {mv_def.mv_name}",
                result.error,
            )

#5.3 事件驱动刷新(EVENT_DRIVEN)

事件驱动刷新在源数据变更时触发:

Python
class EventDrivenRefreshHandler:
    """事件驱动刷新处理器"""

    def __init__(self):
        self._refresh_handler = ManualRefreshHandler()
        self._debounce_windows: dict[str, datetime] = {}
        self._debounce_interval = timedelta(seconds=30)

    async def register(self, mv_def: MaterializedViewDefinition) -> None:
        """注册事件监听"""
        for table in mv_def.source_tables:
            await self._event_bus.subscribe(
                topic=f"data_change:{table}",
                handler=self._on_data_change,
                metadata={"mv_id": mv_def.mv_id},
            )

    async def _on_data_change(
        self, event: DataChangeEvent, metadata: dict
    ) -> None:
        """数据变更事件处理"""
        mv_id = metadata["mv_id"]
        mv_def = await self._metadata_store.get(mv_id)
        if mv_def is None:
            return

        # 防抖:30 秒内的多次变更只触发一次刷新
        last_trigger = self._debounce_windows.get(mv_id)
        now = datetime.utcnow()
        if last_trigger and (now - last_trigger) < self._debounce_interval:
            return

        self._debounce_windows[mv_id] = now

        # 异步触发刷新
        await self._task_queue.enqueue(
            "mv_refresh",
            {"mv_id": mv_id, "trigger": "data_change",
             "source_table": event.table_name},
        )

    async def _process_refresh_task(self, task: dict) -> None:
        mv_def = await self._metadata_store.get(task["mv_id"])
        result = await self._refresh_handler.refresh(mv_def)
        logger.info(
            f"Event-driven refresh for {mv_def.mv_name}: "
            f"{'success' if result.success else 'failed'}, "
            f"trigger={task['trigger']}"
        )

#6. 过期检测(Staleness Detection)

#6.1 过期判定逻辑

Python
class StalenessDetector:
    """物化视图过期检测器"""

    CHECK_INTERVAL = 60  # 每分钟检查一次

    async def check_all(self) -> list[StalenessReport]:
        """检查所有活跃 MV 的新鲜度"""
        all_mvs = await self._metadata_store.list_active()
        reports = []

        for mv in all_mvs:
            report = await self._check_one(mv)
            if report.is_stale:
                reports.append(report)
                await self._handle_stale(mv, report)

        return reports

    async def _check_one(
        self, mv: MaterializedViewDefinition
    ) -> StalenessReport:
        """检查单个 MV 的新鲜度"""
        now = datetime.utcnow()

        # 1. 基于最后刷新时间
        if mv.last_refresh_at is None:
            return StalenessReport(
                mv_name=mv.mv_name,
                is_stale=True,
                reason="never_refreshed",
                age_seconds=None,
            )

        age = (now - mv.last_refresh_at).total_seconds()
        if age > mv.freshness_threshold_seconds:
            return StalenessReport(
                mv_name=mv.mv_name,
                is_stale=True,
                reason="threshold_exceeded",
                age_seconds=age,
                threshold_seconds=mv.freshness_threshold_seconds,
            )

        # 2. 检查源数据是否有更新
        source_updated = await self._check_source_updates(mv)
        if source_updated and age > 60:  # 源数据更新且 MV 超过 1 分钟
            return StalenessReport(
                mv_name=mv.mv_name,
                is_stale=True,
                reason="source_data_updated",
                age_seconds=age,
            )

        return StalenessReport(
            mv_name=mv.mv_name,
            is_stale=False,
            age_seconds=age,
        )

    async def _check_source_updates(
        self, mv: MaterializedViewDefinition
    ) -> bool:
        """检查源表是否有新数据"""
        for table in mv.source_tables:
            last_update = await self._doris.execute(f"""
                SELECT MAX(updated_at) AS last_update
                FROM {table}
            """)
            if last_update and last_update[0]["last_update"]:
                source_time = last_update[0]["last_update"]
                if source_time > mv.last_refresh_at:
                    return True
        return False

    async def _handle_stale(
        self, mv: MaterializedViewDefinition, report: StalenessReport
    ) -> None:
        """处理过期的 MV"""
        # 更新状态
        mv.status = MvStatus.STALE
        await self._metadata_store.save(mv)

        # 发送告警
        await self._alert_service.send(
            AlertLevel.WARNING,
            f"Materialized view {mv.mv_name} is stale",
            f"Age: {report.age_seconds:.0f}s, "
            f"Threshold: {mv.freshness_threshold_seconds}s, "
            f"Reason: {report.reason}",
        )

        # 如果是定时刷新模式,触发立即刷新
        if mv.refresh_mode == RefreshMode.SCHEDULED:
            await self._trigger_immediate_refresh(mv)

#6.2 过期检测流程图

Code
StalenessDetector (runs every 60s)
       |
       v
  List all ACTIVE MVs
       |
       v
  For each MV:
  +-------------------------------------------+
  | last_refresh_at is NULL? --> STALE         |
  |       |NO                                  |
  |       v                                    |
  | age > threshold? --> STALE                 |
  |       |NO                                  |
  |       v                                    |
  | source data updated since refresh?         |
  |       |YES & age > 60s --> STALE           |
  |       |NO                                  |
  |       v                                    |
  | FRESH                                      |
  +-------------------------------------------+
       |
  If STALE:
  ├── Update status to STALE
  ├── Send alert
  └── Trigger immediate refresh (if SCHEDULED)

#7. Schema 变更时的级联失效

#7.1 为什么需要级联失效

当源表的 Schema 发生变更时(增删改列、修改类型),依赖该表的物化视图可能变得无效。如果不处理,MV 查询会返回错误或错误数据。

#7.2 Schema 变更监听

Python
class SchemaChangeListener:
    """Schema 变更监听器"""

    async def on_schema_change(self, event: SchemaChangeEvent) -> None:
        """处理 Schema 变更事件"""
        affected_table = event.table_name
        change_type = event.change_type

        # 查找依赖该表的所有 MV
        affected_mvs = await self._metadata_store.find_by_source_table(
            affected_table
        )

        if not affected_mvs:
            return

        logger.info(
            f"Schema change on {affected_table} ({change_type}): "
            f"affects {len(affected_mvs)} MVs"
        )

        for mv in affected_mvs:
            await self._handle_affected_mv(mv, event)

    async def _handle_affected_mv(
        self, mv: MaterializedViewDefinition, event: SchemaChangeEvent
    ) -> None:
        """处理受影响的 MV"""
        change_type = event.change_type

        if change_type == SchemaChangeType.DROP_COLUMN:
            if self._mv_uses_column(mv, event.column_name):
                # MV 使用了被删除的列 —— 必须重建
                await self._invalidate_and_rebuild(mv, event)
            else:
                # MV 不使用该列,只需刷新
                await self._trigger_refresh(mv)

        elif change_type == SchemaChangeType.ALTER_COLUMN_TYPE:
            if self._mv_uses_column(mv, event.column_name):
                await self._invalidate_and_rebuild(mv, event)

        elif change_type == SchemaChangeType.ADD_COLUMN:
            # 新增列不影响现有 MV
            pass

        elif change_type == SchemaChangeType.RENAME_COLUMN:
            if self._mv_uses_column(mv, event.old_column_name):
                await self._invalidate_and_rebuild(mv, event)

        elif change_type == SchemaChangeType.DROP_TABLE:
            await self._drop_mv(mv, reason="source table dropped")

    async def _invalidate_and_rebuild(
        self, mv: MaterializedViewDefinition, event: SchemaChangeEvent
    ) -> None:
        """失效并重建 MV"""
        # 1. 标记为重建中
        mv.status = MvStatus.REBUILDING
        await self._metadata_store.save(mv)

        # 2. 删除旧 MV
        await self._doris.execute(
            f"DROP MATERIALIZED VIEW IF EXISTS {mv.mv_name}"
        )

        # 3. 重新加载指标定义
        metric = await self._metric_registry.get(mv.metric_id)
        if metric is None:
            mv.status = MvStatus.ERROR
            mv.error_message = "Associated metric not found"
            await self._metadata_store.save(mv)
            return

        # 4. 重新创建 MV
        try:
            new_mv = await self._mv_creator.create_for_metric(metric)
            logger.info(
                f"Successfully rebuilt MV {mv.mv_name} "
                f"after schema change: {event.change_type}"
            )
        except MvCreationError as e:
            mv.status = MvStatus.ERROR
            mv.error_message = f"Rebuild failed: {e}"
            await self._metadata_store.save(mv)
            await self._alert_service.send(
                AlertLevel.ERROR,
                f"MV rebuild failed: {mv.mv_name}",
                str(e),
            )

    def _mv_uses_column(
        self, mv: MaterializedViewDefinition, column_name: str
    ) -> bool:
        """检查 MV 的查询是否使用了指定列"""
        return column_name.lower() in mv.query_sql.lower()

#7.3 Schema 变更类型与 MV 影响矩阵

Schema 变更MV 使用该列MV 不使用该列
ADD_COLUMN无影响无影响
DROP_COLUMN失效 + 重建仅刷新
ALTER_COLUMN_TYPE失效 + 重建无影响
RENAME_COLUMN失效 + 重建无影响
DROP_TABLE删除 MV删除 MV
TRUNCATE_TABLE刷新刷新

#8. MV 健康监控

Python
class MvHealthMonitor:
    """物化视图健康监控"""

    async def get_health_report(self) -> MvHealthReport:
        """获取所有 MV 的健康报告"""
        all_mvs = await self._metadata_store.list_all()

        total = len(all_mvs)
        by_status = {}
        stale_mvs = []
        error_mvs = []
        slow_refreshes = []

        for mv in all_mvs:
            status = mv.status.value
            by_status[status] = by_status.get(status, 0) + 1

            if mv.status == MvStatus.STALE:
                stale_mvs.append(mv.mv_name)
            elif mv.status == MvStatus.ERROR:
                error_mvs.append({
                    "name": mv.mv_name,
                    "error": mv.error_message,
                })

            if mv.last_refresh_duration_ms > 60000:
                slow_refreshes.append({
                    "name": mv.mv_name,
                    "duration_ms": mv.last_refresh_duration_ms,
                })

        return MvHealthReport(
            total_mvs=total,
            status_distribution=by_status,
            stale_mvs=stale_mvs,
            error_mvs=error_mvs,
            slow_refreshes=slow_refreshes,
            overall_health="HEALTHY" if not error_mvs else "DEGRADED",
        )

#9. gRPC 服务定义

PROTOBUF
service MaterializedViewService {
  rpc CreateMV(CreateMVRequest) returns (MaterializedViewDefinition);
  rpc DropMV(DropMVRequest) returns (google.protobuf.Empty);
  rpc RefreshMV(RefreshMVRequest) returns (RefreshResult);
  rpc GetMVStatus(GetMVStatusRequest) returns (MaterializedViewDefinition);
  rpc ListMVs(ListMVsRequest) returns (ListMVsResponse);
  rpc GetHealthReport(google.protobuf.Empty) returns (MvHealthReport);
  rpc UpdateRefreshSchedule(UpdateRefreshScheduleRequest)
      returns (MaterializedViewDefinition);
}

#10. Doris 物化视图特性深度利用

#10.1 异步物化视图

Doris 2.1+ 支持异步物化视图,这是智策平台 MV 自动化的技术基础:

SQL
-- 创建异步物化视图
CREATE MATERIALIZED VIEW mv_customer_order_count
BUILD DEFERRED
REFRESH ASYNC START('2024-01-01 00:00:00') EVERY(INTERVAL 1 HOUR)
DISTRIBUTED BY HASH(object_id) BUCKETS AUTO
AS
SELECT
    c.object_id,
    COUNT(o.object_id) AS metric_value,
    NOW() AS computed_at
FROM customer_objects c
LEFT JOIN order_objects o ON o.customer_id = c.object_id
GROUP BY c.object_id;

-- 手动触发刷新
REFRESH MATERIALIZED VIEW mv_customer_order_count;

-- 查看 MV 状态
SHOW CREATE MATERIALIZED VIEW mv_customer_order_count;

#10.2 查询自动路由

Doris 的查询优化器可以自动识别查询是否命中物化视图,无需用户显式查询 MV 表:

SQL
-- 用户查询(不知道 MV 存在)
SELECT customer_id, COUNT(*) AS order_count
FROM order_objects
GROUP BY customer_id;

-- Doris 自动路由到 mv_customer_order_count(如果匹配)
-- EXPLAIN 会显示: MaterializedView: mv_customer_order_count

#11. 端到端流程

Code
Step 1: Register Metric
  MetricRegistryService.register({
    name: "customer_order_count",
    strategies: [MATERIALIZED, CACHED, REALTIME],
    freshness: HOURLY
  })
       |
       v
Step 2: Auto-Create MV
  MaterializedViewService.create_for_metric()
  ├── Generate SQL: CREATE MATERIALIZED VIEW mv_customer_order_count ...
  ├── Execute on Doris
  ├── Register hourly refresh: "0 * * * *"
  └── Initial refresh
       |
       v
Step 3: Query Uses MV
  OQL: SELECT metric('customer_order_count') FROM Customer
  OQL Rewriter → JOIN mv_customer_order_count
  Doris → fast index scan (< 20ms)
       |
       v
Step 4: Staleness Detection (runs every 60s)
  StalenessDetector → check age vs threshold
  If stale → alert + trigger refresh
       |
       v
Step 5: Schema Change
  ALTER TABLE order_objects DROP COLUMN some_column
       |
       v
  SchemaChangeListener.on_schema_change()
  If MV uses dropped column:
  ├── Drop old MV
  ├── Regenerate SQL from metric definition
  └── Create new MV + initial refresh

#Key Takeaways

  1. 注册即创建 — 指标注册时自动生成 CREATE MV SQL、执行创建、注册刷新调度
  2. 三种刷新模式匹配不同场景 — 手动适合低频、定时适合批量、事件驱动适合近实时
  3. 过期检测是主动的 — 不等用户查询发现过期,而是定时巡检 + 源数据变更检测
  4. Schema 变更自动级联 — 列删除/类型变更/重命名都会触发 MV 重建,用户无感知
  5. Doris 异步 MV 是技术基础 — 利用 Doris 原生异步物化视图能力,减少自建逻辑
  6. 健康监控闭环 — 状态跟踪、告警、自动修复形成完整闭环

#Next Article

下一篇 S3-16 实体 360° 视图:一个 API 返回对象的所有维度 将展示如何将属性、关系、指标、历史、事件、血缘和可用操作聚合为统一的实体视图。

Tags: #MaterializedView #AutoMV #RefreshScheduling #StalenessDetection #SchemaChange #Doris #MetricSystem #OntologyPlatform