返回博客

策略路由模式:ComputationCoordinator 的 7 级优先级调度

在一个复杂的本体驱动决策平台中,计算请求的种类五花八门:

Coomia发布于 2025年12月16日14 分钟阅读
分享本文Twitter / X

策略路由模式:ComputationCoordinator 的 7 级优先级调度

系列:S10 设计模式 · 第 2 篇 | 难度:高级 | 阅读时间:18 分钟

#TL;DR

  • 策略路由模式(Strategy Routing Pattern)将请求分发逻辑从硬编码 if-else 提升为可配置、可扩展的优先级路由链。
  • coomia-dip 的 ComputationCoordinator 实现了 7 级优先级调度:从紧急实时决策到低优先级批量分析,每级有独立的资源配额、超时策略和降级方案。
  • 该模式使系统在极端负载下仍能保证高优先级计算的 SLA,同时充分利用空闲资源处理低优先级任务。

#引言:当一个请求不知道该去哪里

在一个复杂的本体驱动决策平台中,计算请求的种类五花八门:

Code
实时风险评估     → 必须 50ms 内返回
仪表板刷新       → 可接受 500ms
推理链执行       → 可能需要 10 秒
批量属性派生     → 可以排队等几分钟
数据回填         → 夜间跑就行

传统做法是在 API 网关或入口层写一堆 if-else:

Python
if request.type == "realtime_risk":
    route_to(fast_pool)
elif request.type == "dashboard":
    route_to(medium_pool)
elif request.type == "batch":
    route_to(slow_pool)
# ... 每加一种类型就加一个分支

这种方式的问题显而易见:路由逻辑和业务类型强耦合,每次新增计算类型都要改代码、重新部署。更糟糕的是,当系统过载时,所有请求一视同仁地被拒绝——关键决策和可选分析同归于尽。

策略路由模式的核心思想是:将路由决策建模为一个独立的、可插拔的策略链,由优先级、资源状态和业务上下文共同决定请求的去向

#一、策略路由模式的定义

#1.1 模式结构

策略路由模式由四个核心组件构成:

Code
┌──────────────────────────────────────────────────────┐
│                  ComputationCoordinator              │
│                                                      │
│  ┌─────────────┐   ┌──────────────┐   ┌───────────┐ │
│  │ Priority     │   │ Resource     │   │ Strategy  │ │
│  │ Classifier   │──▶│ Allocator    │──▶│ Router    │ │
│  └─────────────┘   └──────────────┘   └───────────┘ │
│         │                  │                  │      │
│         ▼                  ▼                  ▼      │
│  ┌─────────────┐   ┌──────────────┐   ┌───────────┐ │
│  │ Context      │   │ Quota        │   │ Fallback  │ │
│  │ Evaluator    │   │ Manager      │   │ Chain     │ │
│  └─────────────┘   └──────────────┘   └───────────┘ │
└──────────────────────────────────────────────────────┘
  • Priority Classifier(优先级分类器):根据请求的元数据、来源和业务上下文,确定其优先级。
  • Resource Allocator(资源分配器):根据当前资源使用情况和配额策略,分配计算资源。
  • Strategy Router(策略路由器):根据优先级和可用资源,选择最合适的执行路径。
  • Context Evaluator(上下文评估器):评估业务上下文(如租户 SLA、操作紧急度)以调整优先级。
  • Quota Manager(配额管理器):管理每个优先级层级的资源配额和令牌桶。
  • Fallback Chain(降级链):当首选路径不可用时,按优先级依次降级。

#1.2 与经典策略模式的区别

经典 GoF 策略模式关注的是"同一接口的不同实现可以互换"。策略路由模式更进一步:

维度经典策略模式策略路由模式
关注点算法互换请求调度
选择依据静态配置动态多因素
降级支持内置降级链
资源感知感知资源配额
优先级多级优先级

#二、ComputationCoordinator 的 7 级优先级体系

#2.1 优先级定义

coomia-dip 的 ComputationCoordinator 定义了 7 个优先级层级,每个层级有明确的 SLA 和资源保证:

Python
class ComputationPriority(IntEnum):
    """7-level priority for computation requests."""

    P0_EMERGENCY = 0       # 紧急决策:实时风控、告警响应
    P1_INTERACTIVE = 1     # 交互式操作:用户触发的 Action
    P2_REALTIME = 2        # 实时分析:仪表板、实时聚合
    P3_NEAR_REALTIME = 3   # 准实时:推理链、规则引擎
    P4_STANDARD = 4        # 标准:派生属性计算、物化视图刷新
    P5_BACKGROUND = 5      # 后台:数据质量检查、索引重建
    P6_BATCH = 6           # 批量:数据回填、历史分析

#2.2 每级的 SLA 契约

YAML
priority_sla:
  P0_EMERGENCY:
    max_latency: 50ms
    resource_guarantee: 20%    # 始终保留 20% 资源
    preemptible: false
    retry_policy: immediate_3x
    fallback: in_memory_cache

  P1_INTERACTIVE:
    max_latency: 200ms
    resource_guarantee: 15%
    preemptible: false
    retry_policy: exponential_3x
    fallback: degraded_result

  P2_REALTIME:
    max_latency: 500ms
    resource_guarantee: 15%
    preemptible: false
    retry_policy: exponential_3x
    fallback: stale_cache

  P3_NEAR_REALTIME:
    max_latency: 5s
    resource_guarantee: 15%
    preemptible: true          # 可被 P0-P2 抢占
    retry_policy: exponential_5x
    fallback: async_queue

  P4_STANDARD:
    max_latency: 30s
    resource_guarantee: 10%
    preemptible: true
    retry_policy: queue_retry
    fallback: delayed_execution

  P5_BACKGROUND:
    max_latency: 5min
    resource_guarantee: 5%
    preemptible: true
    retry_policy: scheduled_retry
    fallback: next_window

  P6_BATCH:
    max_latency: best_effort
    resource_guarantee: 0%     # 仅使用空闲资源
    preemptible: true
    retry_policy: daily_retry
    fallback: skip_and_log

#2.3 优先级分类算法

请求的优先级不是固定的,而是由多个因素动态计算得出的:

Python
class PriorityClassifier:
    """Classify computation requests into priority levels."""

    def classify(self, request: ComputationRequest) -> ComputationPriority:
        base_priority = self._get_base_priority(request.computation_type)

        # Factor 1: 租户 SLA 等级调整
        tenant_adjustment = self._tenant_sla_adjustment(request.tenant_id)

        # Factor 2: 操作上下文调整
        context_adjustment = self._context_adjustment(request.context)

        # Factor 3: 系统负载调整
        load_adjustment = self._load_adjustment()

        # Factor 4: 时间窗口调整(如交易时段提升优先级)
        time_adjustment = self._time_window_adjustment(request.timestamp)

        final_score = (
            base_priority.value
            + tenant_adjustment
            + context_adjustment
            + load_adjustment
            + time_adjustment
        )

        return ComputationPriority(
            max(0, min(6, round(final_score)))
        )

    def _tenant_sla_adjustment(self, tenant_id: str) -> float:
        sla = self.tenant_registry.get_sla(tenant_id)
        return {
            "platinum": -1.0,   # 白金租户提升一级
            "gold": -0.5,
            "silver": 0.0,
            "bronze": 0.5,
        }.get(sla.tier, 0.0)

    def _context_adjustment(self, context: dict) -> float:
        adjustment = 0.0
        if context.get("triggered_by") == "alert":
            adjustment -= 2.0   # 告警触发的计算大幅提升
        if context.get("is_retry"):
            adjustment -= 0.5   # 重试适当提升
        if context.get("user_waiting"):
            adjustment -= 1.0   # 用户在等结果
        return adjustment

#三、资源分配与配额管理

#3.1 令牌桶算法

每个优先级层级使用独立的令牌桶来控制并发度:

Python
class PriorityQuotaManager:
    """Manage resource quotas per priority level using token buckets."""

    def __init__(self, config: QuotaConfig):
        self.buckets: dict[ComputationPriority, TokenBucket] = {}
        for priority in ComputationPriority:
            bucket_config = config.get_bucket_config(priority)
            self.buckets[priority] = TokenBucket(
                capacity=bucket_config.max_concurrent,
                refill_rate=bucket_config.refill_per_second,
                burst_allowance=bucket_config.burst_factor,
            )

    def try_acquire(self, priority: ComputationPriority) -> bool:
        """Try to acquire a token for the given priority."""
        if self.buckets[priority].try_consume(1):
            return True

        # 高优先级可以抢占低优先级的资源
        if priority.value <= 2:  # P0-P2 可抢占
            return self._try_preempt(priority)

        return False

    def _try_preempt(self, priority: ComputationPriority) -> bool:
        """Preempt lower priority tasks to free resources."""
        for lower in reversed(ComputationPriority):
            if lower.value <= priority.value:
                continue
            preempted = self.preemption_manager.preempt_lowest(
                pool=lower, count=1
            )
            if preempted:
                return True
        return False

#3.2 动态配额调整

系统根据实际负载动态调整各级的资源比例:

Python
class DynamicQuotaAdjuster:
    """Adjust quotas based on real-time system load."""

    def adjust(self, metrics: SystemMetrics) -> dict[ComputationPriority, float]:
        """Return adjusted quota percentages."""
        if metrics.cpu_usage > 0.8:
            # 高负载:收缩低优先级,保护高优先级
            return {
                ComputationPriority.P0_EMERGENCY: 0.30,
                ComputationPriority.P1_INTERACTIVE: 0.25,
                ComputationPriority.P2_REALTIME: 0.20,
                ComputationPriority.P3_NEAR_REALTIME: 0.15,
                ComputationPriority.P4_STANDARD: 0.07,
                ComputationPriority.P5_BACKGROUND: 0.03,
                ComputationPriority.P6_BATCH: 0.00,
            }
        elif metrics.cpu_usage < 0.3:
            # 低负载:放宽低优先级
            return {
                ComputationPriority.P0_EMERGENCY: 0.15,
                ComputationPriority.P1_INTERACTIVE: 0.10,
                ComputationPriority.P2_REALTIME: 0.10,
                ComputationPriority.P3_NEAR_REALTIME: 0.15,
                ComputationPriority.P4_STANDARD: 0.15,
                ComputationPriority.P5_BACKGROUND: 0.15,
                ComputationPriority.P6_BATCH: 0.20,
            }
        else:
            return self._default_quotas()

#四、策略路由器的实现

#4.1 路由链架构

策略路由器使用责任链模式,依次评估每个路由策略:

Python
class StrategyRouter:
    """Route computation requests through a chain of strategies."""

    def __init__(self):
        self.strategies: list[RoutingStrategy] = [
            LocalCacheStrategy(),       # 优先命中本地缓存
            DedicatedPoolStrategy(),    # 专用计算池
            SharedPoolStrategy(),       # 共享计算池
            RemoteNodeStrategy(),       # 远程节点
            DegradedModeStrategy(),     # 降级模式
        ]

    async def route(
        self,
        request: ComputationRequest,
        priority: ComputationPriority
    ) -> ComputationResult:
        context = RoutingContext(request=request, priority=priority)

        for strategy in self.strategies:
            if strategy.can_handle(context):
                try:
                    return await strategy.execute(context)
                except StrategyFailure as e:
                    context.add_failure(strategy.name, e)
                    continue  # 尝试下一个策略

        raise AllStrategiesExhausted(
            f"No strategy could handle request {request.id}",
            failures=context.failures,
        )

#4.2 专用池策略

对于高优先级请求,使用预留的专用计算池:

Python
class DedicatedPoolStrategy(RoutingStrategy):
    """Route to dedicated compute pools for high-priority requests."""

    def can_handle(self, context: RoutingContext) -> bool:
        return context.priority.value <= 2  # P0-P2 使用专用池

    async def execute(self, context: RoutingContext) -> ComputationResult:
        pool = self.pool_registry.get_dedicated_pool(context.priority)

        worker = await pool.acquire(
            timeout=context.priority_sla.max_wait_for_worker
        )

        try:
            result = await worker.compute(
                request=context.request,
                deadline=context.priority_sla.max_latency,
            )
            self.metrics.record_latency(
                priority=context.priority,
                latency=result.elapsed,
            )
            return result
        finally:
            pool.release(worker)

#五、降级策略

#5.1 多级降级链

当主路径失败时,系统按优先级选择降级方案:

Python
class FallbackChain:
    """Execute fallback strategies in priority order."""

    FALLBACK_MAP = {
        ComputationPriority.P0_EMERGENCY: [
            InMemoryCacheFallback(),
            LastKnownGoodFallback(),
            ManualOverrideFallback(),
        ],
        ComputationPriority.P1_INTERACTIVE: [
            StaleCacheFallback(max_age=timedelta(seconds=30)),
            DegradedResultFallback(),
            RetryQueueFallback(),
        ],
        ComputationPriority.P2_REALTIME: [
            StaleCacheFallback(max_age=timedelta(minutes=1)),
            ApproximateResultFallback(),
            AsyncNotificationFallback(),
        ],
        # ... 低优先级的降级策略更宽松
    }

    async def execute(
        self,
        priority: ComputationPriority,
        request: ComputationRequest
    ) -> ComputationResult:
        fallbacks = self.FALLBACK_MAP[priority]

        for fb in fallbacks:
            try:
                result = await fb.attempt(request)
                result.metadata["degraded"] = True
                result.metadata["fallback_strategy"] = fb.name
                return result
            except FallbackFailure:
                continue

        raise CriticalFailure(
            f"All fallbacks exhausted for P{priority.value} request"
        )

#5.2 缓存降级

紧急请求的缓存降级策略:

Python
class InMemoryCacheFallback(FallbackStrategy):
    """Return cached result when computation fails."""

    async def attempt(self, request: ComputationRequest) -> ComputationResult:
        cache_key = self._compute_cache_key(request)
        cached = self.cache.get(cache_key)

        if cached is None:
            raise FallbackFailure("No cached result available")

        if cached.age > timedelta(minutes=5):
            self.alerter.warn(
                f"Serving stale cache (age={cached.age}) for {request.id}"
            )

        return ComputationResult(
            value=cached.value,
            source="cache",
            freshness=cached.timestamp,
            confidence=max(0.5, 1.0 - cached.age.total_seconds() / 300),
        )

#六、抢占式调度

#6.1 抢占决策

当高优先级请求到达但资源不足时,系统需要决定是否抢占低优先级任务:

Python
class PreemptionManager:
    """Manage preemption of lower-priority tasks."""

    def should_preempt(
        self,
        incoming: ComputationPriority,
        running: list[RunningTask],
    ) -> list[RunningTask]:
        """Determine which running tasks to preempt."""
        # 规则 1:P0 可以抢占任何非 P0 任务
        # 规则 2:P1 只能抢占 P4 及以下
        # 规则 3:P2 只能抢占 P5 及以下
        # 规则 4:P3 及以下不能抢占

        preemption_threshold = {
            ComputationPriority.P0_EMERGENCY: 1,    # 可抢占 P1+
            ComputationPriority.P1_INTERACTIVE: 4,  # 可抢占 P4+
            ComputationPriority.P2_REALTIME: 5,     # 可抢占 P5+
        }

        threshold = preemption_threshold.get(incoming)
        if threshold is None:
            return []

        candidates = [
            task for task in running
            if task.priority.value >= threshold
        ]

        # 优先抢占已运行时间最长的低优先级任务
        candidates.sort(
            key=lambda t: (-t.priority.value, -t.elapsed.total_seconds())
        )

        return candidates[:1]  # 每次最多抢占一个

    async def preempt(self, task: RunningTask) -> None:
        """Preempt a running task, saving its checkpoint."""
        checkpoint = await task.save_checkpoint()
        await self.checkpoint_store.save(task.id, checkpoint)
        await task.cancel(reason="preempted")

        # 将被抢占的任务重新加入队列
        await self.requeue(task, checkpoint)

#6.2 检查点与恢复

被抢占的任务通过检查点机制恢复执行:

Python
class CheckpointManager:
    """Save and restore computation checkpoints."""

    async def save_checkpoint(self, task: RunningTask) -> Checkpoint:
        return Checkpoint(
            task_id=task.id,
            priority=task.priority,
            progress=task.progress_percentage,
            state=await task.serialize_state(),
            saved_at=datetime.utcnow(),
            preempted_by=task.preempted_by,
        )

    async def restore_and_resume(self, checkpoint: Checkpoint) -> None:
        task = await self.task_factory.create_from_checkpoint(checkpoint)

        # 恢复后优先级提升 0.5 级(避免反复被抢占)
        adjusted_priority = max(
            0, checkpoint.priority.value - 0.5
        )
        task.priority = ComputationPriority(round(adjusted_priority))

        await self.scheduler.enqueue(task)

#七、监控与可观测性

#7.1 关键指标

Python
class RoutingMetrics:
    """Metrics for the strategy routing system."""

    def __init__(self):
        self.request_count = Counter(
            "computation_requests_total",
            "Total computation requests",
            ["priority", "strategy", "status"],
        )
        self.latency = Histogram(
            "computation_latency_seconds",
            "Computation latency",
            ["priority"],
            buckets=[0.01, 0.05, 0.1, 0.5, 1, 5, 30, 300],
        )
        self.preemption_count = Counter(
            "computation_preemptions_total",
            "Total preemptions",
            ["preemptor_priority", "victim_priority"],
        )
        self.fallback_count = Counter(
            "computation_fallbacks_total",
            "Total fallback activations",
            ["priority", "fallback_strategy"],
        )
        self.queue_depth = Gauge(
            "computation_queue_depth",
            "Current queue depth",
            ["priority"],
        )

#7.2 告警规则

YAML
alerts:
  - name: P0LatencyExceeded
    condition: computation_latency_seconds{priority="P0"} > 0.05
    severity: critical
    action: page_oncall

  - name: PreemptionStorm
    condition: rate(computation_preemptions_total[5m]) > 10
    severity: warning
    action: scale_up_pool

  - name: FallbackActivated
    condition: increase(computation_fallbacks_total{priority="P0"}[1m]) > 0
    severity: critical
    action: page_oncall_and_incident

  - name: QueueBacklog
    condition: computation_queue_depth{priority="P4"} > 1000
    severity: warning
    action: notify_channel

#八、与 coomia-dip 架构的集成

#8.1 跨 Layer 路由

ComputationCoordinator 作为 Reasoning & Decision Layer(推理与决策)的核心组件,与其他 Layer 的交互全部通过 gRPC:

PROTOBUF
service ComputationCoordinator {
    rpc SubmitComputation(ComputationRequest) returns (ComputationResponse);
    rpc GetComputationStatus(StatusRequest) returns (StatusResponse);
    rpc CancelComputation(CancelRequest) returns (CancelResponse);
    rpc StreamResults(StreamRequest) returns (stream ComputationResult);
}

message ComputationRequest {
    string request_id = 1;
    string tenant_id = 2;
    string computation_type = 3;
    bytes payload = 4;
    map<string, string> context = 5;
    int32 priority_hint = 6;  // 调用方的优先级建议
}

#8.2 与 Ontology 层的联动

路由策略可以通过 Ontology 模型动态配置:

Python
# 通过 Ontology 定义路由规则
routing_rule = platform.objects.RoutingRule.create(
    name="high_value_customer_boost",
    condition="request.context.customer_tier == 'enterprise'",
    priority_adjustment=-1,  # 提升一级
    effective_from=datetime(2026, 1, 1),
    effective_to=datetime(2026, 12, 31),
)

#九、实战案例:金融风控场景

#9.1 场景描述

一个银行部署的 coomia-dip 实例需要同时处理:

  • 实时交易风控(P0):每笔交易必须在 50ms 内完成风险评估
  • 客户画像更新(P3):准实时更新客户的行为画像
  • 反洗钱分析(P5):后台运行的复杂图分析
  • 月度合规报告(P6):月底批量生成

#9.2 系统行为

Code
09:30 开盘:交易洪峰到来
  → P6 批量任务暂停
  → P5 后台任务降速
  → P0/P1 获得 60% 资源

11:00 交易平稳期
  → 动态配额恢复正常分配
  → P5 恢复正常速度
  → P6 开始处理积压任务

14:55 盘前波动
  → P0 延迟告警触发
  → P3 中两个推理任务被抢占
  → 释放资源后 P0 延迟恢复正常

22:00 夜间维护窗口
  → P6 批量任务全速运行
  → P5 索引重建启动

#十、总结

#Key Takeaways

  1. 7 级优先级不是拍脑袋定的——每一级对应一类真实的业务场景,有明确的 SLA 和降级策略。
  2. 动态分类比静态配置更实用——请求的优先级由租户 SLA、操作上下文、系统负载和时间窗口共同决定。
  3. 抢占式调度是保证 P0 SLA 的关键——没有抢占,高优先级就只是个标签,没有实际保障。
  4. 降级链比硬失败更友好——即使在极端负载下,用户看到的是降级结果而非错误页面。
  5. 可观测性决定模式的成败——没有完善的监控和告警,优先级调度就是黑箱。

#参考资料

  1. Linux CFS Scheduler
  2. Palantir Foundry Computation Service
  3. coomia-dip 架构总览
  4. Gamma et al., Design Patterns, Addison-Wesley, 1994
  5. Google Borg: Large-Scale Cluster Management

tags: strategy-pattern routing priority-scheduling preemption computation coomia-dip

下一篇:S10-03 事件溯源