Back to Blog

S3-15 Materialized View Automation: Register a Metric, Get a Materialized View

The coomia-dip MaterializedViewService implements "register a metric, get a materialized view" automation. It supports three refresh modes (manual / scheduled / event-driven), provides staleness detection, and automatically cascades invalidation on schema changes. This article fully dissects the auto-creation, refresh scheduling, staleness detection, and schema-change cascading pipeline.

CoomiaPublished on July 25, 202512 min read
Share this articleTwitter / X

S3-15 Materialized View Automation: Register a Metric, Get a Materialized View

Series: S3 Data Foundation · Article 15 | Level: Advanced | Reading Time: 20 min

#TL;DR

The coomia-dip MaterializedViewService implements "register a metric, get a materialized view" automation. It supports three refresh modes (manual / scheduled / event-driven), provides staleness detection, and automatically cascades invalidation on schema changes. This article fully dissects the auto-creation, refresh scheduling, staleness detection, and schema-change cascading pipeline.

#1. Why Materialized View Automation Is Needed

Materialized Views (MVs) pre-compute query results and persist them for fast access. In traditional databases, creating and maintaining MVs requires manual DBA work: writing CREATE MATERIALIZED VIEW statements, configuring refresh strategies, monitoring freshness, and rebuilding MVs after schema changes.

In an Ontology platform, metrics are dynamically registered. Operations teams may add or modify metric definitions daily. If every metric's MV requires manual DBA maintenance, this process simply cannot scale.

The coomia-dip design goal: automatically create MVs when metrics are registered, automatically destroy them when metrics are deleted, and automatically rebuild them on schema changes.

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 Architecture

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 Metadata Model

#3.1 MaterializedViewDefinition

Python
class MaterializedViewDefinition(BaseModel):
    """Materialized view definition"""
    mv_id: str = Field(default_factory=lambda: str(uuid4()))
    mv_name: str                       # mv_customer_order_count
    metric_id: str                     # associated metric ID
    metric_name: str                   # associated metric name
    object_type: str                   # Customer
    source_tables: list[str]           # source table list
    create_sql: str                    # CREATE MV statement
    query_sql: str                     # internal query SQL
    refresh_mode: RefreshMode          # refresh mode
    refresh_schedule: str | None = None # cron expression
    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. Automatic MV Creation

#4.1 Generating MVs from Metric Definitions

Python
class MvCreator:
    """Materialized view creator"""

    async def create_for_metric(
        self, metric: MetricDefinition
    ) -> MaterializedViewDefinition:
        """Automatically create MV based on metric definition"""

        mv_name = f"mv_{metric.name}"
        query_sql = self._generate_query_sql(metric)
        create_sql = self._generate_create_sql(mv_name, query_sql, metric)
        refresh_mode = self._determine_refresh_mode(metric)
        refresh_schedule = self._determine_schedule(metric)

        try:
            await self._doris.execute(create_sql)
        except DorisError as e:
            raise MvCreationError(f"Failed to create MV {mv_name}: {e}")

        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)

        if refresh_mode == RefreshMode.SCHEDULED:
            await self._scheduler.register(mv_def)
        elif refresh_mode == RefreshMode.EVENT_DRIVEN:
            await self._event_listener.register(mv_def)

        await self._trigger_refresh(mv_def)
        return mv_def

    def _generate_query_sql(self, metric: MetricDefinition) -> str:
        """Generate query SQL from metric expression"""
        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:
            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:
        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:
        schedules = {
            FreshnessLevel.HOURLY: "0 * * * *",
            FreshnessLevel.DAILY: "0 2 * * *",
            FreshnessLevel.WEEKLY: "0 2 * * 1",
        }
        return schedules.get(metric.freshness_requirement)

#5. Three Refresh Modes

#5.1 Manual Refresh

Python
class ManualRefreshHandler:
    """Manual refresh handler"""

    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 Refresh

Python
class ScheduledRefreshHandler:
    """Scheduled refresh handler"""

    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 _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 Refresh

Event-driven refresh triggers when source data changes:

Python
class EventDrivenRefreshHandler:
    """Event-driven refresh handler"""

    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

        # Debounce: multiple changes within 30s trigger only one refresh
        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},
        )

#6. Staleness Detection

#6.1 Staleness Logic

Python
class StalenessDetector:
    """Materialized view staleness detector"""

    CHECK_INTERVAL = 60  # check every minute

    async def check_all(self) -> list[StalenessReport]:
        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:
        now = datetime.utcnow()

        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,
            )

        source_updated = await self._check_source_updates(mv)
        if source_updated and age > 60:
            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 _handle_stale(
        self, mv: MaterializedViewDefinition, report: StalenessReport
    ) -> None:
        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 Staleness Detection Flow

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                                      |
  +-------------------------------------------+

#7. Schema Change Cascading Invalidation

#7.1 Why Cascading Invalidation Is Needed

When source table schemas change (add/drop/alter columns, change types), dependent MVs may become invalid. Without handling this, MV queries return errors or incorrect data.

#7.2 Schema Change Listener

Python
class SchemaChangeListener:
    """Schema change listener"""

    async def on_schema_change(self, event: SchemaChangeEvent) -> None:
        affected_table = event.table_name

        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} ({event.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:
        change_type = event.change_type

        if change_type == SchemaChangeType.DROP_COLUMN:
            if self._mv_uses_column(mv, event.column_name):
                await self._invalidate_and_rebuild(mv, event)
            else:
                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:
            pass  # new columns don't affect existing MVs

        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.status = MvStatus.REBUILDING
        await self._metadata_store.save(mv)

        await self._doris.execute(
            f"DROP MATERIALIZED VIEW IF EXISTS {mv.mv_name}"
        )

        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

        try:
            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)

    def _mv_uses_column(
        self, mv: MaterializedViewDefinition, column_name: str
    ) -> bool:
        return column_name.lower() in mv.query_sql.lower()

#7.3 Schema Change vs MV Impact Matrix

Schema ChangeMV Uses ColumnMV Does Not Use Column
ADD_COLUMNNo impactNo impact
DROP_COLUMNInvalidate + RebuildRefresh only
ALTER_COLUMN_TYPEInvalidate + RebuildNo impact
RENAME_COLUMNInvalidate + RebuildNo impact
DROP_TABLEDrop MVDrop MV
TRUNCATE_TABLERefreshRefresh

#8. MV Health Monitoring

Python
class MvHealthMonitor:
    """Materialized view health monitor"""

    async def get_health_report(self) -> MvHealthReport:
        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 Service Definition

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 Materialized View Deep Integration

#10.1 Async Materialized Views

Doris 2.1+ supports async materialized views, which form the technical foundation for coomia-dip MV automation:

SQL
-- Create async materialized view
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;

-- Manual refresh trigger
REFRESH MATERIALIZED VIEW mv_customer_order_count;

-- View MV status
SHOW CREATE MATERIALIZED VIEW mv_customer_order_count;

#10.2 Automatic Query Routing

Doris query optimizer can automatically detect whether a query hits a materialized view without the user explicitly querying the MV table:

SQL
-- User query (unaware of MV)
SELECT customer_id, COUNT(*) AS order_count
FROM order_objects
GROUP BY customer_id;

-- Doris automatically routes to mv_customer_order_count (if matching)
-- EXPLAIN shows: MaterializedView: mv_customer_order_count

#11. End-to-End Flow

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 ...
  +-- 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. Register equals create — metric registration automatically generates CREATE MV SQL, executes creation, and registers refresh scheduling
  2. Three refresh modes match different scenarios — manual for low-frequency, scheduled for batch, event-driven for near-real-time
  3. Staleness detection is proactive — does not wait for user queries to discover staleness; periodic scans + source data change detection
  4. Schema changes cascade automatically — column drops, type changes, and renames trigger MV rebuilds transparently
  5. Doris async MVs are the technical foundation — leverages Doris native async materialized view capabilities to minimize custom logic
  6. Health monitoring closes the loop — status tracking, alerting, and auto-repair form a complete feedback loop

#Next Article

Next up: S3-16 Entity 360 View: One API Call Returns All Object Dimensions will show how to aggregate attributes, relations, metrics, history, events, lineage, and available actions into a unified entity view.

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