返回博客

级联模式:依赖传播与影响分析

在本体驱动的系统中,对象之间通过 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

  1. Schema 声明式:级联策略在 LinkType Schema 中声明,由平台自动执行
  2. 四种策略:CASCADE、RESTRICT、SET_NULL、SET_DEFAULT 覆盖所有常见场景
  3. 拓扑排序:依赖 DAG 和拓扑排序确保级联操作按正确顺序执行
  4. 影响分析:变更前的影响评估让用户了解操作的波及范围和风险等级
  5. 异步执行:大规模级联通过异步工作流和批量处理避免超时
  6. 派生属性:源属性变更自动触发派生属性重新计算,形成完整的级联链

#Next Article

下一篇我们将探讨查询改写模式(Query Rewrite)——coomia-dip 如何将高层语义查询改写为底层存储引擎的优化查询。

S10-09: 查询改写:从语义到存储的优化翻译

#Tags

#设计模式 #级联模式 #Cascade #依赖传播 #影响分析 #拓扑排序 #DAG #派生属性 #级联删除