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.
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.
+------------------------------------------------------------------+
| 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
+------------------------------------------------------------------+
| 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
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
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
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
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:
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
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
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
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 Change | MV Uses Column | MV Does Not Use Column |
|---|---|---|
| ADD_COLUMN | No impact | No impact |
| DROP_COLUMN | Invalidate + Rebuild | Refresh only |
| ALTER_COLUMN_TYPE | Invalidate + Rebuild | No impact |
| RENAME_COLUMN | Invalidate + Rebuild | No impact |
| DROP_TABLE | Drop MV | Drop MV |
| TRUNCATE_TABLE | Refresh | Refresh |
#8. MV Health Monitoring
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
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:
-- 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:
-- 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
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
- Register equals create — metric registration automatically generates CREATE MV SQL, executes creation, and registers refresh scheduling
- Three refresh modes match different scenarios — manual for low-frequency, scheduled for batch, event-driven for near-real-time
- Staleness detection is proactive — does not wait for user queries to discover staleness; periodic scans + source data change detection
- Schema changes cascade automatically — column drops, type changes, and renames trigger MV rebuilds transparently
- Doris async MVs are the technical foundation — leverages Doris native async materialized view capabilities to minimize custom logic
- 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