策略路由模式:ComputationCoordinator 的 7 级优先级调度
在一个复杂的本体驱动决策平台中,计算请求的种类五花八门:
策略路由模式:ComputationCoordinator 的 7 级优先级调度
“系列:S10 设计模式 · 第 2 篇 | 难度:高级 | 阅读时间:18 分钟
#TL;DR
- 策略路由模式(Strategy Routing Pattern)将请求分发逻辑从硬编码 if-else 提升为可配置、可扩展的优先级路由链。
- coomia-dip 的 ComputationCoordinator 实现了 7 级优先级调度:从紧急实时决策到低优先级批量分析,每级有独立的资源配额、超时策略和降级方案。
- 该模式使系统在极端负载下仍能保证高优先级计算的 SLA,同时充分利用空闲资源处理低优先级任务。
#引言:当一个请求不知道该去哪里
在一个复杂的本体驱动决策平台中,计算请求的种类五花八门:
实时风险评估 → 必须 50ms 内返回
仪表板刷新 → 可接受 500ms
推理链执行 → 可能需要 10 秒
批量属性派生 → 可以排队等几分钟
数据回填 → 夜间跑就行
传统做法是在 API 网关或入口层写一堆 if-else:
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 模式结构
策略路由模式由四个核心组件构成:
┌──────────────────────────────────────────────────────┐
│ 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 和资源保证:
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 契约
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 优先级分类算法
请求的优先级不是固定的,而是由多个因素动态计算得出的:
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 令牌桶算法
每个优先级层级使用独立的令牌桶来控制并发度:
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 动态配额调整
系统根据实际负载动态调整各级的资源比例:
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 路由链架构
策略路由器使用责任链模式,依次评估每个路由策略:
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 专用池策略
对于高优先级请求,使用预留的专用计算池:
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 多级降级链
当主路径失败时,系统按优先级选择降级方案:
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 缓存降级
紧急请求的缓存降级策略:
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 抢占决策
当高优先级请求到达但资源不足时,系统需要决定是否抢占低优先级任务:
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 检查点与恢复
被抢占的任务通过检查点机制恢复执行:
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 关键指标
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 告警规则
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:
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 模型动态配置:
# 通过 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 系统行为
09:30 开盘:交易洪峰到来
→ P6 批量任务暂停
→ P5 后台任务降速
→ P0/P1 获得 60% 资源
11:00 交易平稳期
→ 动态配额恢复正常分配
→ P5 恢复正常速度
→ P6 开始处理积压任务
14:55 盘前波动
→ P0 延迟告警触发
→ P3 中两个推理任务被抢占
→ 释放资源后 P0 延迟恢复正常
22:00 夜间维护窗口
→ P6 批量任务全速运行
→ P5 索引重建启动
#十、总结
#Key Takeaways
- 7 级优先级不是拍脑袋定的——每一级对应一类真实的业务场景,有明确的 SLA 和降级策略。
- 动态分类比静态配置更实用——请求的优先级由租户 SLA、操作上下文、系统负载和时间窗口共同决定。
- 抢占式调度是保证 P0 SLA 的关键——没有抢占,高优先级就只是个标签,没有实际保障。
- 降级链比硬失败更友好——即使在极端负载下,用户看到的是降级结果而非错误页面。
- 可观测性决定模式的成败——没有完善的监控和告警,优先级调度就是黑箱。
#参考资料
- Linux CFS Scheduler↗
- Palantir Foundry Computation Service↗
- coomia-dip 架构总览
- Gamma et al., Design Patterns, Addison-Wesley, 1994
- Google Borg: Large-Scale Cluster Management↗
tags: strategy-pattern routing priority-scheduling preemption computation coomia-dip
下一篇:S10-03 事件溯源