级联模式:依赖传播与影响分析
在本体驱动的系统中,对象之间通过 LinkType 建立丰富的关联关系。修改一个对象可能触发连锁反应:
Coomia发布于 2025年12月23日12 分钟阅读
分享本文Twitter / X
级联模式:依赖传播与影响分析
“系列:S10 设计模式 · 第 8 篇 | 难度:高级 | 阅读时间:18 分钟
#TL;DR
- 级联模式(Cascade Pattern)定义了 Ontology 对象之间依赖变更如何传播——当一个对象变更时,所有依赖它的对象应如何响应。
- coomia-dip 支持四种级联策略:级联更新(Cascade Update)、级联删除(Cascade Delete)、限制(Restrict)和置空(Set Null),并在 LinkType 定义中声明。
- 通过依赖 DAG(有向无环图)和拓扑排序,coomia-dip 确保级联操作按正确的顺序执行,并提供变更前的影响分析(Impact Analysis)能力。
#引言:蝴蝶效应
在本体驱动的系统中,对象之间通过 LinkType 建立丰富的关联关系。修改一个对象可能触发连锁反应:
Code
场景:删除一个 Department(部门)
→ 该部门下的 Employee(员工)怎么办?
→ 员工名下的 Asset(资产)怎么办?
→ 资产关联的 MaintenanceRecord(维修记录)怎么办?
→ 员工参与的 Project(项目)怎么办?
→ 项目关联的 Budget(预算)怎么办?
如果没有明确的级联策略,开发者要么手动处理每一层依赖(容易遗漏),要么忽略依赖导致数据不一致(悬挂引用)。
级联模式将这些策略声明在 Schema 中,由平台自动执行。
#一、级联策略定义
#1.1 四种级联策略
Python
from dataclasses import dataclass, field
from enum import Enum
from typing import Any
class CascadeStrategy(Enum):
"""Cascade strategies for dependency propagation."""
CASCADE = "cascade" # 级联操作:删除父对象时同时删除子对象
RESTRICT = "restrict" # 限制操作:如果有依赖对象则拒绝操作
SET_NULL = "set_null" # 置空:将依赖引用设为 null
SET_DEFAULT = "set_default" # 设默认值:将依赖引用设为默认值
NO_ACTION = "no_action" # 不操作:不做任何处理(危险)
@dataclass
class CascadeRule:
"""Cascade rule defined on a LinkType."""
link_type: str
source_type: str
target_type: str
on_delete: CascadeStrategy = CascadeStrategy.RESTRICT
on_update: CascadeStrategy = CascadeStrategy.CASCADE
on_archive: CascadeStrategy = CascadeStrategy.CASCADE
depth_limit: int = 10 # 最大级联深度
async_execution: bool = False # 是否异步执行级联
batch_size: int = 1000 # 批量处理大小
@dataclass
class CascadeConfig:
"""Complete cascade configuration for an ObjectType."""
object_type: str
rules: list[CascadeRule] = field(default_factory=list)
def get_rules_for_source(self, source_type: str) -> list[CascadeRule]:
"""Get cascade rules where this type is the source."""
return [r for r in self.rules if r.source_type == source_type]
#1.2 在 LinkType 中声明级联策略
Python
# 定义部门-员工关系的级联策略
department_employee_link = CascadeRule(
link_type="Department_employs_Employee",
source_type="Department",
target_type="Employee",
on_delete=CascadeStrategy.RESTRICT, # 有员工时不允许删除部门
on_update=CascadeStrategy.CASCADE, # 部门更新时级联更新员工的部门引用
on_archive=CascadeStrategy.CASCADE, # 部门归档时级联归档员工
)
# 定义员工-资产关系
employee_asset_link = CascadeRule(
link_type="Employee_owns_Asset",
source_type="Employee",
target_type="Asset",
on_delete=CascadeStrategy.SET_NULL, # 员工删除时,资产的 owner 设为 null
on_update=CascadeStrategy.CASCADE,
)
# 定义项目-预算关系
project_budget_link = CascadeRule(
link_type="Project_has_Budget",
source_type="Project",
target_type="Budget",
on_delete=CascadeStrategy.CASCADE, # 项目删除时级联删除预算
on_update=CascadeStrategy.CASCADE,
)
#二、依赖图与拓扑排序
#2.1 构建依赖 DAG
coomia-dip 在启动时根据 LinkType 和 CascadeRule 构建全局的依赖有向无环图(DAG):
Python
class DependencyDAG:
"""Directed Acyclic Graph for object dependencies."""
def __init__(self):
self._edges: dict[str, list[tuple[str, CascadeRule]]] = {}
self._reverse_edges: dict[str, list[tuple[str, CascadeRule]]] = {}
def add_dependency(self, rule: CascadeRule) -> None:
"""Add a dependency edge to the DAG."""
self._edges.setdefault(rule.source_type, []).append(
(rule.target_type, rule)
)
self._reverse_edges.setdefault(rule.target_type, []).append(
(rule.source_type, rule)
)
def get_dependents(self, object_type: str) -> list[tuple[str, CascadeRule]]:
"""Get all types that depend on the given type."""
return self._edges.get(object_type, [])
def get_dependencies(self, object_type: str) -> list[tuple[str, CascadeRule]]:
"""Get all types that the given type depends on."""
return self._reverse_edges.get(object_type, [])
def topological_sort(self) -> list[str]:
"""Return types in topological order (parents before children)."""
in_degree: dict[str, int] = {}
all_types: set[str] = set()
for source, targets in self._edges.items():
all_types.add(source)
for target, _ in targets:
all_types.add(target)
in_degree[target] = in_degree.get(target, 0) + 1
queue = [t for t in all_types if in_degree.get(t, 0) == 0]
result = []
while queue:
node = queue.pop(0)
result.append(node)
for target, _ in self._edges.get(node, []):
in_degree[target] -= 1
if in_degree[target] == 0:
queue.append(target)
if len(result) != len(all_types):
raise CyclicDependencyError(
"Circular dependency detected in Ontology schema"
)
return result
def get_cascade_order(self, object_type: str) -> list[str]:
"""Get the order in which cascade operations should be executed."""
visited: set[str] = set()
order: list[str] = []
def dfs(current: str, depth: int = 0) -> None:
if current in visited or depth > 20:
return
visited.add(current)
for target, _ in self.get_dependents(current):
dfs(target, depth + 1)
order.append(current)
dfs(object_type)
return list(reversed(order))
#2.2 循环依赖检测
coomia-dip 在 Schema 注册时自动检测循环依赖,防止无限级联:
Python
class CyclicDependencyDetector:
"""Detect circular dependencies in the Ontology schema."""
def detect(self, dag: DependencyDAG) -> list[list[str]]:
"""Find all cycles in the dependency graph."""
cycles: list[list[str]] = []
visited: set[str] = set()
rec_stack: set[str] = set()
def dfs(node: str, path: list[str]) -> None:
visited.add(node)
rec_stack.add(node)
path.append(node)
for target, _ in dag.get_dependents(node):
if target not in visited:
dfs(target, path)
elif target in rec_stack:
# 找到循环
cycle_start = path.index(target)
cycles.append(path[cycle_start:] + [target])
path.pop()
rec_stack.discard(node)
all_types = set()
for source in dag._edges:
all_types.add(source)
for target, _ in dag._edges[source]:
all_types.add(target)
for node in all_types:
if node not in visited:
dfs(node, [])
return cycles
#三、级联执行引擎
#3.1 级联删除
Python
class CascadeEngine:
"""Engine for executing cascade operations."""
def __init__(
self,
dag: DependencyDAG,
repository: "ObjectRepository",
event_store: "EventStore",
):
self._dag = dag
self._repo = repository
self._event_store = event_store
async def cascade_delete(
self,
object_type: str,
object_id: str,
context: "OperationContext",
dry_run: bool = False,
) -> "CascadeResult":
"""Execute cascade delete with impact analysis."""
# Phase 1: 影响分析
impact = await self._analyze_impact(
object_type, object_id, "delete"
)
if dry_run:
return CascadeResult(
status="dry_run",
impact=impact,
executed=False,
)
# Phase 2: 检查 RESTRICT 策略
for dep_type, rule in self._dag.get_dependents(object_type):
if rule.on_delete == CascadeStrategy.RESTRICT:
count = await self._repo.count_by_reference(
dep_type, object_type, object_id
)
if count > 0:
return CascadeResult(
status="restricted",
impact=impact,
executed=False,
error=f"Cannot delete: {count} {dep_type} objects depend on this",
)
# Phase 3: 按拓扑逆序执行级联(先删子对象)
cascade_order = self._dag.get_cascade_order(object_type)
deleted_objects: list[dict] = []
for dep_type in reversed(cascade_order):
if dep_type == object_type:
continue
rule = self._get_rule(object_type, dep_type)
if not rule:
continue
dependents = await self._repo.find_by_reference(
dep_type, object_type, object_id
)
for dep_obj in dependents:
if rule.on_delete == CascadeStrategy.CASCADE:
# 递归级联删除
await self.cascade_delete(
dep_type, dep_obj["_id"], context
)
deleted_objects.append(dep_obj)
elif rule.on_delete == CascadeStrategy.SET_NULL:
await self._repo.set_reference_null(
dep_type, dep_obj["_id"], object_type
)
elif rule.on_delete == CascadeStrategy.SET_DEFAULT:
default = rule.metadata.get("default_value")
await self._repo.set_reference(
dep_type, dep_obj["_id"], object_type, default
)
# Phase 4: 删除目标对象
await self._repo.delete(object_type, object_id)
# Phase 5: 记录事件
await self._event_store.append([
DomainEvent(
event_id=generate_id(),
event_type="cascade.delete",
aggregate_id=object_id,
aggregate_type=object_type,
sequence_number=0,
timestamp=datetime.utcnow(),
payload={
"cascade_deleted": len(deleted_objects),
"impact": impact,
},
metadata=EventMetadata(
actor_id=context.actor_id,
actor_type=context.actor_type,
tenant_id=context.tenant_id,
world_id=context.world_id,
source_plane="control",
trace_id=context.trace_id,
),
)
])
return CascadeResult(
status="completed",
impact=impact,
executed=True,
deleted_count=len(deleted_objects) + 1,
)
def _get_rule(self, source_type: str, target_type: str) -> CascadeRule | None:
"""Get cascade rule between two types."""
for dep_type, rule in self._dag.get_dependents(source_type):
if dep_type == target_type:
return rule
return None
#3.2 级联更新
Python
class CascadeUpdateEngine:
"""Handle cascade updates when object properties change."""
async def cascade_update(
self,
object_type: str,
object_id: str,
changes: dict[str, Any],
context: "OperationContext",
) -> "CascadeResult":
"""Cascade property changes to dependent objects."""
# 确定哪些变更需要级联
cascadable_changes = self._filter_cascadable(object_type, changes)
if not cascadable_changes:
return CascadeResult(status="no_cascade_needed", executed=False)
updated_objects: list[dict] = []
for dep_type, rule in self._dag.get_dependents(object_type):
if rule.on_update != CascadeStrategy.CASCADE:
continue
dependents = await self._repo.find_by_reference(
dep_type, object_type, object_id
)
for dep_obj in dependents:
# 计算依赖对象需要更新的字段
dep_changes = self._compute_dependent_changes(
rule, changes, dep_obj
)
if dep_changes:
await self._repo.update(dep_type, dep_obj["_id"], dep_changes)
updated_objects.append({
"type": dep_type,
"id": dep_obj["_id"],
"changes": dep_changes,
})
# 递归级联
await self.cascade_update(
dep_type, dep_obj["_id"], dep_changes, context
)
return CascadeResult(
status="completed",
executed=True,
updated_count=len(updated_objects),
impact={"updated_objects": updated_objects},
)
def _filter_cascadable(
self, object_type: str, changes: dict
) -> dict:
"""Filter changes that need to be cascaded."""
# 只有被其他对象引用的字段变更才需要级联
cascadable = {}
for field_name, new_value in changes.items():
if self._is_referenced_field(object_type, field_name):
cascadable[field_name] = new_value
return cascadable
#四、影响分析
#4.1 变更前影响评估
在执行级联操作之前,coomia-dip 提供影响分析能力,让用户了解操作将影响多少对象:
Python
class ImpactAnalyzer:
"""Analyze the impact of cascade operations before execution."""
async def analyze(
self,
object_type: str,
object_id: str,
operation: str,
) -> dict:
"""Analyze the full impact of a cascade operation."""
impact: dict[str, Any] = {
"root": {"type": object_type, "id": object_id},
"operation": operation,
"affected_types": {},
"total_affected": 0,
"risk_level": "low",
}
await self._collect_impact(
object_type, object_id, operation, impact, depth=0
)
# 计算风险等级
if impact["total_affected"] > 1000:
impact["risk_level"] = "critical"
elif impact["total_affected"] > 100:
impact["risk_level"] = "high"
elif impact["total_affected"] > 10:
impact["risk_level"] = "medium"
return impact
async def _collect_impact(
self,
object_type: str,
object_id: str,
operation: str,
impact: dict,
depth: int,
) -> None:
"""Recursively collect impact information."""
if depth > 10:
return
for dep_type, rule in self._dag.get_dependents(object_type):
strategy = getattr(rule, f"on_{operation}", CascadeStrategy.NO_ACTION)
count = await self._repo.count_by_reference(
dep_type, object_type, object_id
)
if count > 0:
impact["affected_types"][dep_type] = {
"count": count,
"strategy": strategy.value,
"depth": depth + 1,
}
impact["total_affected"] += count
if strategy == CascadeStrategy.CASCADE:
# 递归分析子依赖
dependents = await self._repo.find_by_reference(
dep_type, object_type, object_id
)
for dep in dependents[:10]: # 限制递归采样
await self._collect_impact(
dep_type, dep["_id"], operation, impact, depth + 1
)
async def generate_report(
self, impact: dict
) -> str:
"""Generate a human-readable impact report."""
lines = [
f"Impact Analysis Report",
f"=" * 40,
f"Operation: {impact['operation']} on {impact['root']['type']}#{impact['root']['id']}",
f"Risk Level: {impact['risk_level']}",
f"Total Affected Objects: {impact['total_affected']}",
f"",
f"Affected Types:",
]
for type_name, info in impact["affected_types"].items():
lines.append(
f" - {type_name}: {info['count']} objects "
f"(strategy: {info['strategy']}, depth: {info['depth']})"
)
return "\n".join(lines)
#4.2 SDK 中的影响分析 API
Python
# 通过 SDK 使用影响分析
from ontology_sdk import OntoPlatform
client = OntoPlatform.connect("https://platform.example.com")
# 删除前先分析影响
impact = client.objects.analyze_delete_impact(
object_type="Department",
object_id="dept-001",
)
print(f"Risk level: {impact.risk_level}")
print(f"Total affected: {impact.total_affected}")
for type_name, info in impact.affected_types.items():
print(f" {type_name}: {info.count} objects ({info.strategy})")
# 确认后执行删除
if impact.risk_level in ("low", "medium"):
result = client.objects.delete(
object_type="Department",
object_id="dept-001",
cascade=True,
)
#五、大规模级联的异步执行
#5.1 异步级联任务
当级联操作影响大量对象时,同步执行可能导致超时。coomia-dip 支持异步级联:
Python
class AsyncCascadeExecutor:
"""Execute large cascade operations asynchronously."""
async def submit(
self,
object_type: str,
object_id: str,
operation: str,
context: "OperationContext",
) -> str:
"""Submit an async cascade operation."""
task_id = generate_id()
# 提交到 Temporal 工作流
await self._temporal.start_workflow(
CascadeWorkflow,
args={
"task_id": task_id,
"object_type": object_type,
"object_id": object_id,
"operation": operation,
"context": context,
},
id=f"cascade-{task_id}",
)
return task_id
async def get_progress(self, task_id: str) -> dict:
"""Get progress of an async cascade operation."""
return await self._progress_store.get(task_id)
#5.2 批量级联处理
Python
class BatchCascadeProcessor:
"""Process cascade operations in batches."""
async def process_batch(
self,
dep_type: str,
dep_objects: list[dict],
rule: CascadeRule,
operation: str,
batch_size: int = 1000,
) -> dict:
"""Process cascade in batches to avoid memory issues."""
total = len(dep_objects)
processed = 0
errors = []
for i in range(0, total, batch_size):
batch = dep_objects[i:i + batch_size]
for obj in batch:
try:
if operation == "delete" and rule.on_delete == CascadeStrategy.CASCADE:
await self._repo.delete(dep_type, obj["_id"])
elif operation == "delete" and rule.on_delete == CascadeStrategy.SET_NULL:
await self._repo.set_reference_null(
dep_type, obj["_id"], rule.source_type
)
processed += 1
except Exception as e:
errors.append({"id": obj["_id"], "error": str(e)})
# 报告进度
await self._progress_reporter.report(
processed=processed, total=total, errors=len(errors)
)
return {
"total": total,
"processed": processed,
"errors": errors,
}
#六、Derived Property 的级联重算
#6.1 派生属性依赖
当源属性变更时,依赖它的派生属性(Derived Property)需要重新计算。这是一种特殊的级联更新:
Python
class DerivedPropertyCascade:
"""Cascade recalculation for derived properties."""
async def on_property_change(
self,
object_type: str,
object_id: str,
changed_property: str,
new_value: Any,
) -> list[dict]:
"""Trigger recalculation of derived properties when source changes."""
# 查找所有依赖此属性的派生属性
affected = self._dependency_graph.get_derived_from(
object_type, changed_property
)
recalculated = []
for derived in affected:
# 重新计算派生值
new_derived_value = await self._calculator.compute(
derived.formula,
object_type=object_type,
object_id=object_id,
)
await self._repo.update(
derived.target_type,
object_id,
{derived.property_name: new_derived_value},
)
recalculated.append({
"property": derived.property_name,
"old_value": derived.current_value,
"new_value": new_derived_value,
})
# 递归:此派生属性的变更可能触发更多级联
await self.on_property_change(
derived.target_type,
object_id,
derived.property_name,
new_derived_value,
)
return recalculated
#Key Takeaways
- Schema 声明式:级联策略在 LinkType Schema 中声明,由平台自动执行
- 四种策略:CASCADE、RESTRICT、SET_NULL、SET_DEFAULT 覆盖所有常见场景
- 拓扑排序:依赖 DAG 和拓扑排序确保级联操作按正确顺序执行
- 影响分析:变更前的影响评估让用户了解操作的波及范围和风险等级
- 异步执行:大规模级联通过异步工作流和批量处理避免超时
- 派生属性:源属性变更自动触发派生属性重新计算,形成完整的级联链
#Next Article
下一篇我们将探讨查询改写模式(Query Rewrite)——coomia-dip 如何将高层语义查询改写为底层存储引擎的优化查询。
S10-09: 查询改写:从语义到存储的优化翻译
#Tags
#设计模式 #级联模式 #Cascade #依赖传播 #影响分析 #拓扑排序 #DAG #派生属性 #级联删除