计算策略路由:7 种计算模式的优先级调度
TL;DR
计算策略路由:7 种计算模式的优先级调度
“系列:S2 架构全景 · 第 7 篇 | 难度:中级 | 阅读时间:18 分钟
TL;DR
- coomia-dip 平台支持 7 种计算模式(缓存查询、物化视图、SQL 计算、表达式求值、Reducer 聚合、函数计算、实时流计算),ComputationCoordinator 负责在每次属性查询时选择最优的执行路径。
- 路由决策基于成本模型,综合考虑数据新鲜度要求、查询延迟预算、资源消耗和缓存命中率,在毫秒级内完成路径选择。
- 这套计算路由机制是 coomia-dip 派生属性(Derived Property)系统的核心引擎。它让上层业务无需关心数据是从缓存读取、从 Doris 计算还是通过实时流获取——只需声明"我要什么",路由器自动决定"怎么算"。
#1. 引言:为什么需要计算路由
在前一篇关于事件驱动架构的讨论中,我们了解了 Kafka 如何串联平台各组件。但事件触发之后,具体的属性值如何计算?这就是计算路由要解决的问题。
考虑一个简单的场景:一个供应链平台需要计算"供应商综合风险评分"。这个派生属性依赖于:
- 历史交货准时率(需要聚合几百万条交货记录)
- 最近 30 天质量投诉数(需要时间窗口过滤)
- 财务健康指数(需要调用外部 API)
- 同行业排名百分位(需要跨供应商比较)
在传统架构中,每次查询都执行全量计算显然不现实。但如果简单地缓存结果,又可能返回过时数据。coomia-dip 的 ComputationCoordinator 就是为解决这类问题而设计的——它为每次查询动态选择最合适的计算路径。
#2. 七种计算模式概览
#2.1 计算模式分类
coomia-dip 将所有属性计算归纳为 7 种模式,按优先级从高到低排列:
| 优先级 | 模式 | 名称 | 典型延迟 | 适用场景 |
|---|---|---|---|---|
| 1 | CACHE | 缓存查询 | < 1ms | 热点数据、高频读取 |
| 2 | MATERIALIZED | 物化视图 | 1-5ms | 预计算聚合、报表指标 |
| 3 | SQL | SQL 计算 | 5-50ms | 复杂查询、多表关联 |
| 4 | EXPRESSION | 表达式求值 | 1-10ms | 简单公式、字段组合 |
| 5 | REDUCER | Reducer 聚合 | 10-100ms | 大规模聚合、窗口计算 |
| 6 | FUNCTION | 函数计算 | 50-500ms | 自定义逻辑、外部调用 |
| 7 | STREAMING | 实时流计算 | 100ms-5s | 实时聚合、CEP 规则 |
#2.2 优先级不等于执行顺序
一个关键概念:优先级决定的是"偏好"而非"固定顺序"。ComputationCoordinator 会根据多个因素动态调整实际选择:
路由决策 = f(数据新鲜度要求, 延迟预算, 缓存状态, 资源可用性, 数据规模)
例如,即使缓存优先级最高,但如果业务要求 "新鲜度 < 1 秒" 而缓存已经 30 秒未刷新,路由器会跳过缓存直接选择 SQL 或 STREAMING 模式。
#3. ComputationCoordinator 核心架构
#3.1 组件定位
ComputationCoordinator 位于 Control Layer (Control Layer) 的 OntologyRuntimeService 内部,是派生属性查询的核心调度器。
┌─────────────────────────────────────────────────┐
│ OntologyRuntimeService │
│ │
│ ┌─────────────────────────────────────────┐ │
│ │ ComputationCoordinator │ │
│ │ │ │
│ │ ┌───────────┐ ┌──────────────────┐ │ │
│ │ │ CostModel │ │ FreshnessChecker │ │ │
│ │ └───────────┘ └──────────────────┘ │ │
│ │ ┌───────────┐ ┌──────────────────┐ │ │
│ │ │ RouteTable│ │ CircuitBreaker │ │ │
│ │ └───────────┘ └──────────────────┘ │ │
│ └─────────────────────────────────────────┘ │
│ │
│ ┌──────┐ ┌────────────┐ ┌─────┐ ┌──────────┐ │
│ │Cache │ │Materialized│ │ SQL │ │Expression│ │
│ │Engine│ │ Engine │ │Exec │ │ Evaluator│ │
│ └──────┘ └────────────┘ └─────┘ └──────────┘ │
│ ┌──────────┐ ┌──────────┐ ┌─────────────────┐│
│ │ Reducer │ │ Function │ │ StreamingEngine ││
│ │ Engine │ │ Runtime │ │ ││
│ └──────────┘ └──────────┘ └─────────────────┘│
└─────────────────────────────────────────────────┘
#3.2 路由决策流程
每次派生属性查询都经过以下步骤:
步骤 1: 解析请求
├─ 提取 ObjectType + PropertyName
├─ 读取属性的 ComputationSpec(计算规则定义)
└─ 确定 freshness_requirement 和 latency_budget
步骤 2: 评估候选模式
├─ 检查每种模式的可用性(引擎是否健康)
├─ 检查缓存命中状态和新鲜度
├─ 估算每种模式的执行成本
└─ 过滤不满足约束的模式
步骤 3: 选择最优路径
├─ 对候选模式按成本排序
├─ 选择成本最低且满足约束的模式
└─ 记录路由决策(用于审计和优化)
步骤 4: 执行计算
├─ 将请求分发到选中的计算引擎
├─ 设置超时和熔断保护
└─ 返回结果并异步更新缓存
#3.3 成本模型详解
成本模型是路由决策的核心。每种计算模式的成本由以下公式估算:
def estimate_cost(mode, context):
"""计算模式成本估算"""
base_cost = MODE_BASE_COSTS[mode] # 基础开销
data_cost = context.data_size * MODE_DATA_FACTORS[mode] # 数据规模因子
resource_cost = get_resource_pressure(mode) # 当前资源压力
freshness_penalty = calc_freshness_gap( # 新鲜度惩罚
mode, context.freshness_requirement
)
total = base_cost + data_cost + resource_cost + freshness_penalty
# 缓存命中时成本极低
if mode == CACHE and cache_hit(context):
total *= 0.01
return total
#4. 模式一:缓存查询 (CACHE)
#4.1 缓存层架构
缓存查询是最快的计算路径。coomia-dip 使用两级缓存架构:
L1 缓存: 进程内缓存 (Caffeine)
├─ 容量: 每个服务实例 10,000 条
├─ 过期策略: 写后 30 秒自动过期
└─ 命中延迟: < 0.1ms
L2 缓存: 分布式缓存 (Redis Cluster)
├─ 容量: 无硬性限制,按内存配额管理
├─ 过期策略: 按属性定义的 cache_ttl
└─ 命中延迟: 0.5-2ms
#4.2 缓存失效策略
缓存失效采用"事件驱动 + TTL 双保险"模式:
-
事件驱动失效:当源数据变更时,CDC 事件通过 Kafka 触发缓存失效。从数据变更到缓存失效的延迟通常在 50-200ms。
-
TTL 兜底:即使事件丢失,TTL 保证缓存最终会过期。默认 TTL 根据属性类型设定:
- 静态属性(如名称、编号):TTL = 5 分钟
- 准实时属性(如库存数量):TTL = 30 秒
- 实时属性(如价格):TTL = 5 秒
#4.3 缓存预热
系统启动时,ComputationCoordinator 会根据历史访问频率预热 Top-1000 的热点属性。预热过程异步执行,不阻塞服务启动。
#5. 模式二:物化视图 (MATERIALIZED)
#5.1 物化视图的定位
物化视图适用于计算成本高但变更频率低的聚合类属性。例如:
- "部门年度总营收"——需要聚合数万条订单,但每天只需更新一次
- "产品评价平均分"——需要聚合所有评价,但每小时更新即可
#5.2 Doris 物化视图
coomia-dip 利用 Apache Doris 的物化视图能力实现预计算:
-- 创建物化视图(由 ComputationCoordinator 自动管理)
CREATE MATERIALIZED VIEW mv_supplier_risk_score AS
SELECT
supplier_id,
AVG(delivery_on_time_rate) as avg_delivery_rate,
COUNT(quality_complaint_id) as complaint_count,
MAX(last_assessment_date) as latest_assessment
FROM supplier_delivery_records
LEFT JOIN quality_complaints USING (supplier_id)
GROUP BY supplier_id;
#5.3 物化刷新策略
物化视图的刷新由 ComputationCoordinator 统一调度:
| 刷新模式 | 触发条件 | 适用场景 |
|---|---|---|
| 定时刷新 | Cron 表达式 | 日报、周报类指标 |
| 事件触发 | CDC 事件累积达阈值 | 准实时聚合 |
| 手动刷新 | API 调用 | 临时分析需求 |
| 惰性刷新 | 查询时发现过期 | 低频访问的属性 |
#6. 模式三:SQL 计算
#6.1 动态 SQL 生成
对于无法预计算的属性查询,ComputationCoordinator 会根据属性的 ComputationSpec 动态生成 SQL:
class SQLComputationEngine:
def compute(self, spec: ComputationSpec, context: QueryContext):
# 从 ComputationSpec 构建 SQL
sql_builder = SQLBuilder(spec.object_type)
# 添加计算逻辑
for rule in spec.computation_rules:
sql_builder.add_computation(rule)
# 添加过滤条件
sql_builder.add_filters(context.filters)
# 添加安全边界(防止全表扫描)
sql_builder.add_limit(spec.max_rows or 10000)
sql_builder.add_timeout(context.latency_budget)
# 执行查询
return self.doris_client.execute(sql_builder.build())
#6.2 查询安全保护
SQL 计算模式内置多重安全保护:
- 行数限制:默认最多返回 10,000 行,防止内存溢出
- 超时限制:查询超时默认 30 秒,超时自动终止
- 资源配额:每个租户的并发 SQL 查询数限制为 20
- SQL 注入防护:所有参数通过预编译语句传递,禁止字符串拼接
#7. 模式四:表达式求值 (EXPRESSION)
#7.1 表达式引擎
表达式求值适用于简单的属性组合计算,例如:
// 全名 = 姓 + " " + 名
full_name = last_name + " " + first_name
// 毛利率 = (收入 - 成本) / 收入 * 100
gross_margin = (revenue - cost) / revenue * 100
// 状态标签 = IF(score > 80, "优秀", IF(score > 60, "合格", "不合格"))
status_label = IF(score > 80, "优秀", IF(score > 60, "合格", "不合格"))
#7.2 表达式编译与缓存
为了避免每次执行都解析表达式字符串,引擎会将表达式编译为 AST 并缓存:
源表达式 ──解析──> AST ──优化──> 优化后 AST ──编译──> 可执行函数
│
缓存到内存
(Key = 表达式哈希)
编译后的表达式执行延迟通常在微秒级别,使得表达式模式成为仅次于缓存的最快计算路径。
#7.3 类型安全
表达式引擎在编译阶段执行类型检查:
revenue (Decimal) - cost (Decimal) → Decimal ✓
name (String) + age (Integer) → 类型错误 ✗
类型不兼容的表达式会在属性定义保存时就被拒绝,而不是在运行时抛出异常。
#8. 模式五:Reducer 聚合
#8.1 MapReduce 式分布式聚合
当数据规模超过单机处理能力时,ComputationCoordinator 会选择 Reducer 模式。这种模式将聚合任务分解为 Map 和 Reduce 两个阶段:
Map 阶段 (并行):
Partition 1 ──> 局部聚合结果 1
Partition 2 ──> 局部聚合结果 2
Partition 3 ──> 局部聚合结果 3
...
Partition N ──> 局部聚合结果 N
Reduce 阶段 (汇总):
局部结果 1..N ──> 全局聚合结果
#8.2 支持的聚合操作
| 操作 | 描述 | 是否可并行化 |
|---|---|---|
| SUM | 求和 | 是 |
| COUNT | 计数 | 是 |
| AVG | 平均值 | 是(转为 SUM/COUNT) |
| MIN/MAX | 最值 | 是 |
| DISTINCT_COUNT | 去重计数 | 是(HyperLogLog) |
| PERCENTILE | 百分位数 | 近似可并行(T-Digest) |
| TOP_K | 前 K 项 | 是(合并排序) |
#8.3 Reducer 与 SQL 的选择边界
ComputationCoordinator 根据数据规模自动选择 SQL 或 Reducer:
数据量 < 100 万行 → SQL 模式(单节点 Doris 足够)
数据量 100 万 - 1 亿行 → Reducer 模式(分布式聚合)
数据量 > 1 亿行 → 物化视图(预计算 + 增量更新)
#9. 模式六:函数计算 (FUNCTION)
#9.1 自定义函数注册
函数计算模式允许开发者注册自定义计算逻辑,适用于无法用 SQL 或表达式表达的复杂计算:
# 在 Reasoning & Decision Layer (Intelligence Layer) 注册自定义函数
@computation_function(
name="supplier_risk_score",
input_types={"supplier_id": "string"},
output_type="decimal",
timeout_ms=5000,
cacheable=True,
cache_ttl_seconds=3600
)
async def compute_supplier_risk(supplier_id: str) -> Decimal:
"""综合计算供应商风险评分"""
# 获取多维度数据
delivery_data = await get_delivery_history(supplier_id)
financial_data = await get_financial_health(supplier_id)
complaint_data = await get_quality_complaints(supplier_id)
# 多维度加权计算
score = (
delivery_data.on_time_rate * 0.35 +
financial_data.health_index * 0.30 +
(1 - complaint_data.rate) * 0.20 +
delivery_data.response_speed * 0.15
)
return Decimal(str(round(score, 2)))
#9.2 函数调用链路
函数计算的调用经过 gRPC 跨 Layer 通信:
Control Layer (B) Intelligence Layer (D)
ComputationCoordinator FunctionRuntime
│ │
├──gRPC──> ExecuteFunction ──────────>│
│ (supplier_risk_score, │
│ {supplier_id: "S001"}) │
│ │
│ 计算执行中...
│ │
│<──gRPC── FunctionResult <──────────│
│ (score: 0.78, │
│ compute_time_ms: 230) │
#9.3 函数计算的保护机制
- 超时控制:每个函数有独立的超时设置(默认 5 秒)
- 资源隔离:函数在独立的沙箱中执行,不影响主进程
- 重试策略:幂等函数支持自动重试(最多 3 次)
- 熔断保护:连续失败 5 次后熔断,60 秒后半开尝试
#10. 模式七:实时流计算 (STREAMING)
#10.1 流计算引擎
实时流计算是延迟最高但实时性最强的计算模式。它适用于:
- 需要实时聚合的指标(如"最近 5 分钟订单数")
- 复杂事件处理(CEP)规则(如"连续 3 次异常告警")
- 实时风控评分
#10.2 与 Kafka Streams 的集成
coomia-dip 的流计算引擎构建在 Kafka Streams 之上:
Kafka Topic Streaming Engine 结果
(source events) (output)
│ │
├─> Window(5min) ─> Count ─> "最近5分钟订单数" ────────>│
├─> Window(1h) ─> Avg ─> "小时均价" ────────────────>│
├─> CEP Rule ─> Match ─> "异常模式检测" ───────────>│
└─> Tumbling(1d)─> Sum ─> "日累计金额" ──────────────>│
#10.3 窗口类型
| 窗口类型 | 描述 | 示例 |
|---|---|---|
| Tumbling | 固定窗口,不重叠 | 每小时统计一次 |
| Sliding | 滑动窗口,可重叠 | 最近 5 分钟的移动平均 |
| Session | 会话窗口,按活动间隔 | 用户会话内的操作计数 |
| Global | 全局窗口,无边界 | 累计总数 |
#10.4 流计算的回退策略
当流计算引擎不可用时,ComputationCoordinator 自动回退到 SQL 模式:
async def compute_with_fallback(spec, context):
try:
# 优先尝试流计算
result = await streaming_engine.compute(spec, context)
return result
except StreamingUnavailableError:
# 回退到 SQL 计算(可能延迟更高但仍可用)
logger.warning(
"Streaming engine unavailable, falling back to SQL",
property=spec.property_name
)
return await sql_engine.compute(spec, context)
#11. 路由决策的高级特性
#11.1 多属性批量路由
当一次查询请求多个派生属性时,ComputationCoordinator 会进行批量优化:
请求: 获取供应商 S001 的 [风险评分, 交货准时率, 质量评级, 活跃状态]
批量路由优化:
风险评分 → FUNCTION (需要自定义计算)
交货准时率 → CACHE (高频访问,缓存命中)
质量评级 → MATERIALIZED (预计算物化视图)
活跃状态 → EXPRESSION (简单布尔表达式)
并行执行: CACHE + EXPRESSION 立即返回
MATERIALIZED + FUNCTION 并行计算
合并结果: 最终延迟 = max(各路径延迟) ≈ 230ms
#11.2 自适应路由学习
ComputationCoordinator 会记录每次路由决策的实际执行结果,并据此调整成本模型参数:
historical_stats = {
"supplier_risk_score": {
"CACHE": {"avg_latency": 0.5, "hit_rate": 0.85},
"SQL": {"avg_latency": 45, "success_rate": 0.99},
"FUNCTION": {"avg_latency": 230, "success_rate": 0.97},
"MATERIALIZED": {"avg_latency": 3, "freshness_gap": 3600}
}
}
# 基于历史数据动态调整成本因子
cost_factor["SQL"] = base_cost * (1 + failure_rate * penalty)
#11.3 租户级路由策略
不同租户可以配置不同的路由偏好:
tenant_routing_config:
tenant_gold:
# 金牌租户:优先实时性
freshness_weight: 0.8
latency_weight: 0.1
cost_weight: 0.1
max_streaming_partitions: 16
tenant_silver:
# 银牌租户:均衡策略
freshness_weight: 0.4
latency_weight: 0.3
cost_weight: 0.3
max_streaming_partitions: 8
tenant_basic:
# 基础租户:优先成本
freshness_weight: 0.2
latency_weight: 0.2
cost_weight: 0.6
streaming_enabled: false
#12. 熔断与降级策略
#12.1 分级熔断
ComputationCoordinator 为每种计算引擎独立维护熔断状态:
熔断器状态:
CACHE → CLOSED (正常) ── 连续失败 0/5
MATERIALIZED → CLOSED (正常) ── 连续失败 0/5
SQL → HALF_OPEN (试探) ── 上次熔断 30s 前
EXPRESSION → CLOSED (正常) ── 连续失败 0/5
REDUCER → CLOSED (正常) ── 连续失败 1/5
FUNCTION → OPEN (熔断中) ── 将在 60s 后试探
STREAMING → CLOSED (正常) ── 连续失败 0/5
#12.2 优雅降级路径
当高优先级模式不可用时,路由器自动选择下一最佳模式:
理想路径: CACHE → 命中 → 返回 (0.5ms)
降级路径 1: CACHE → 未命中 → MATERIALIZED → 返回 (3ms)
降级路径 2: CACHE → 未命中 → MATERIALIZED → 过期 → SQL → 返回 (45ms)
降级路径 3: 全部不可用 → 返回默认值 + 告警 (0ms + alert)
#12.3 背压控制
当系统负载过高时,ComputationCoordinator 会主动限制高成本计算模式:
class BackpressureController:
def should_allow(self, mode: ComputationMode) -> bool:
current_load = get_system_load()
if current_load > 0.9: # 系统负载 > 90%
# 只允许 CACHE 和 EXPRESSION
return mode in (CACHE, EXPRESSION)
if current_load > 0.7: # 系统负载 > 70%
# 禁止 STREAMING 和 FUNCTION
return mode not in (STREAMING, FUNCTION)
return True # 正常模式,全部允许
#13. 与 Palantir Foundry 的对比
| 维度 | Palantir Foundry | coomia-dip |
|---|---|---|
| 计算引擎 | Spark (批处理为主) | 7 种模式混合调度 |
| 路由策略 | 固定管道 (Pipeline) | 动态成本路由 |
| 缓存层 | 内置 Object 缓存 | 两级缓存 (Caffeine + Redis) |
| 实时计算 | Foundry Streaming (受限) | Kafka Streams (完整) |
| 物化视图 | Dataset Materialization | Doris 物化视图 |
| 自定义函数 | TypeScript Functions | Python 函数计算 |
| 多租户路由 | 资源隔离 (Namespace) | 策略级路由隔离 |
| 降级策略 | 手动配置 | 自动熔断 + 降级 |
coomia-dip 的核心优势在于动态路由——Foundry 的计算路径在 Pipeline 构建时确定,而 coomia-dip 在每次查询时实时选择最优路径。这意味着面对数据增长、负载变化、组件故障等场景,coomia-dip 能自动调整而无需人工干预。
#14. 可观测性与调优
#14.1 路由决策日志
每次路由决策都会记录详细日志,便于分析和优化:
{
"trace_id": "abc-123",
"property": "supplier_risk_score",
"object_id": "S001",
"candidates": [
{"mode": "CACHE", "cost": 0.5, "available": true, "fresh": false},
{"mode": "MATERIALIZED", "cost": 3.2, "available": true, "fresh": true},
{"mode": "SQL", "cost": 45.0, "available": true, "fresh": true},
{"mode": "FUNCTION", "cost": 230.0, "available": false, "reason": "circuit_open"}
],
"selected": "MATERIALIZED",
"reason": "lowest_cost_meeting_freshness",
"actual_latency_ms": 2.8,
"timestamp": "2026-03-24T10:15:30Z"
}
#14.2 关键监控指标
| 指标 | 描述 | 告警阈值 |
|---|---|---|
| route_cache_hit_rate | 缓存命中率 | < 60% 告警 |
| route_fallback_rate | 降级比率 | > 20% 告警 |
| route_decision_latency_p99 | 路由决策延迟 P99 | > 5ms 告警 |
| compute_timeout_rate | 计算超时率 | > 5% 告警 |
| circuit_breaker_open_count | 熔断器打开数 | > 2 告警 |
#14.3 性能调优建议
- 提高缓存命中率:分析 cache miss 日志,对高频属性调整 TTL 或增加预热范围
- 物化视图覆盖:对 SQL 模式的高频查询,考虑创建物化视图
- 函数计算优化:对 > 200ms 的函数,检查是否可以拆分为缓存 + 增量计算
- 减少降级:熔断频繁的引擎需要排查根因(通常是资源不足或网络问题)
#Key Takeaways
-
7 种计算模式覆盖了从微秒到秒级的完整延迟光谱。 ComputationCoordinator 不是简单地"尝试缓存,失败就查数据库",而是基于成本模型对 7 种模式进行全局最优选择。每种模式有明确的适用场景、性能特征和保护机制,形成完整的计算策略矩阵。
-
动态路由让平台具备自适应能力。 传统系统的计算路径在开发时确定,一旦数据量增长或负载变化就需要人工调整。coomia-dip 的成本模型基于实时指标和历史统计动态调整,路由决策自动适应系统状态变化。结合多租户策略配置,不同业务场景可以获得定制化的计算体验。
-
熔断降级保证系统永不因计算引擎故障而完全不可用。 7 种计算模式本身就是一个天然的降级链。即使最坏情况下只剩下表达式引擎可用,系统仍然能返回基本计算结果。这种分级降级策略让平台的可用性从单点可用提升到了矩阵可用。
“下一篇预告: [S2-08] API 设计哲学:Ontology-Native API 的设计原则——深入理解 coomia-dip 如何将本体概念映射为 RESTful API 和 gRPC 接口,实现"以本体为中心"的 API 体验。
Tags: #computation-routing #derived-property #cost-model #cache #materialized-view #sql #expression #reducer #function #streaming #circuit-breaker #coomia-dip #智策平台