返回博客

计算策略路由:7 种计算模式的优先级调度

TL;DR

Coomia发布于 2025年6月30日20 分钟阅读
分享本文Twitter / X

计算策略路由: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 种模式,按优先级从高到低排列:

优先级模式名称典型延迟适用场景
1CACHE缓存查询< 1ms热点数据、高频读取
2MATERIALIZED物化视图1-5ms预计算聚合、报表指标
3SQLSQL 计算5-50ms复杂查询、多表关联
4EXPRESSION表达式求值1-10ms简单公式、字段组合
5REDUCERReducer 聚合10-100ms大规模聚合、窗口计算
6FUNCTION函数计算50-500ms自定义逻辑、外部调用
7STREAMING实时流计算100ms-5s实时聚合、CEP 规则

#2.2 优先级不等于执行顺序

一个关键概念:优先级决定的是"偏好"而非"固定顺序"。ComputationCoordinator 会根据多个因素动态调整实际选择:

Code
路由决策 = f(数据新鲜度要求, 延迟预算, 缓存状态, 资源可用性, 数据规模)

例如,即使缓存优先级最高,但如果业务要求 "新鲜度 < 1 秒" 而缓存已经 30 秒未刷新,路由器会跳过缓存直接选择 SQL 或 STREAMING 模式。

#3. ComputationCoordinator 核心架构

#3.1 组件定位

ComputationCoordinator 位于 Control Layer (Control Layer) 的 OntologyRuntimeService 内部,是派生属性查询的核心调度器。

Code
┌─────────────────────────────────────────────────┐
│              OntologyRuntimeService             │
│                                                 │
│  ┌─────────────────────────────────────────┐    │
│  │       ComputationCoordinator            │    │
│  │                                         │    │
│  │  ┌───────────┐  ┌──────────────────┐   │    │
│  │  │ CostModel │  │ FreshnessChecker │   │    │
│  │  └───────────┘  └──────────────────┘   │    │
│  │  ┌───────────┐  ┌──────────────────┐   │    │
│  │  │ RouteTable│  │ CircuitBreaker   │   │    │
│  │  └───────────┘  └──────────────────┘   │    │
│  └─────────────────────────────────────────┘    │
│                                                 │
│  ┌──────┐ ┌────────────┐ ┌─────┐ ┌──────────┐ │
│  │Cache │ │Materialized│ │ SQL │ │Expression│ │
│  │Engine│ │  Engine    │ │Exec │ │ Evaluator│ │
│  └──────┘ └────────────┘ └─────┘ └──────────┘ │
│  ┌──────────┐ ┌──────────┐ ┌─────────────────┐│
│  │ Reducer  │ │ Function │ │ StreamingEngine ││
│  │ Engine   │ │ Runtime  │ │                 ││
│  └──────────┘ └──────────┘ └─────────────────┘│
└─────────────────────────────────────────────────┘

#3.2 路由决策流程

每次派生属性查询都经过以下步骤:

Code
步骤 1: 解析请求
  ├─ 提取 ObjectType + PropertyName
  ├─ 读取属性的 ComputationSpec(计算规则定义)
  └─ 确定 freshness_requirement 和 latency_budget

步骤 2: 评估候选模式
  ├─ 检查每种模式的可用性(引擎是否健康)
  ├─ 检查缓存命中状态和新鲜度
  ├─ 估算每种模式的执行成本
  └─ 过滤不满足约束的模式

步骤 3: 选择最优路径
  ├─ 对候选模式按成本排序
  ├─ 选择成本最低且满足约束的模式
  └─ 记录路由决策(用于审计和优化)

步骤 4: 执行计算
  ├─ 将请求分发到选中的计算引擎
  ├─ 设置超时和熔断保护
  └─ 返回结果并异步更新缓存

#3.3 成本模型详解

成本模型是路由决策的核心。每种计算模式的成本由以下公式估算:

Python
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 使用两级缓存架构:

Code
L1 缓存: 进程内缓存 (Caffeine)
  ├─ 容量: 每个服务实例 10,000 条
  ├─ 过期策略: 写后 30 秒自动过期
  └─ 命中延迟: < 0.1ms

L2 缓存: 分布式缓存 (Redis Cluster)
  ├─ 容量: 无硬性限制,按内存配额管理
  ├─ 过期策略: 按属性定义的 cache_ttl
  └─ 命中延迟: 0.5-2ms

#4.2 缓存失效策略

缓存失效采用"事件驱动 + TTL 双保险"模式:

  1. 事件驱动失效:当源数据变更时,CDC 事件通过 Kafka 触发缓存失效。从数据变更到缓存失效的延迟通常在 50-200ms。

  2. 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 的物化视图能力实现预计算:

SQL
-- 创建物化视图(由 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:

Python
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 表达式引擎

表达式求值适用于简单的属性组合计算,例如:

Code
// 全名 = 姓 + " " + 名
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 并缓存:

Code
源表达式 ──解析──> AST ──优化──> 优化后 AST ──编译──> 可执行函数
                                                        │
                                                     缓存到内存
                                                     (Key = 表达式哈希)

编译后的表达式执行延迟通常在微秒级别,使得表达式模式成为仅次于缓存的最快计算路径。

#7.3 类型安全

表达式引擎在编译阶段执行类型检查:

Code
revenue (Decimal) - cost (Decimal) → Decimal ✓
name (String) + age (Integer) → 类型错误 ✗

类型不兼容的表达式会在属性定义保存时就被拒绝,而不是在运行时抛出异常。

#8. 模式五:Reducer 聚合

#8.1 MapReduce 式分布式聚合

当数据规模超过单机处理能力时,ComputationCoordinator 会选择 Reducer 模式。这种模式将聚合任务分解为 Map 和 Reduce 两个阶段:

Code
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:

Code
数据量 < 100 万行 → SQL 模式(单节点 Doris 足够)
数据量 100 万 - 1 亿行 → Reducer 模式(分布式聚合)
数据量 > 1 亿行 → 物化视图(预计算 + 增量更新)

#9. 模式六:函数计算 (FUNCTION)

#9.1 自定义函数注册

函数计算模式允许开发者注册自定义计算逻辑,适用于无法用 SQL 或表达式表达的复杂计算:

Python
# 在 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 通信:

Code
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 之上:

Code
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 模式:

Python
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 会进行批量优化:

Code
请求: 获取供应商 S001 的 [风险评分, 交货准时率, 质量评级, 活跃状态]

批量路由优化:
  风险评分     → FUNCTION (需要自定义计算)
  交货准时率   → CACHE (高频访问,缓存命中)
  质量评级     → MATERIALIZED (预计算物化视图)
  活跃状态     → EXPRESSION (简单布尔表达式)

并行执行: CACHE + EXPRESSION 立即返回
         MATERIALIZED + FUNCTION 并行计算
合并结果: 最终延迟 = max(各路径延迟) ≈ 230ms

#11.2 自适应路由学习

ComputationCoordinator 会记录每次路由决策的实际执行结果,并据此调整成本模型参数:

Code
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 租户级路由策略

不同租户可以配置不同的路由偏好:

YAML
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 为每种计算引擎独立维护熔断状态:

Code
熔断器状态:
  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 优雅降级路径

当高优先级模式不可用时,路由器自动选择下一最佳模式:

Code
理想路径:    CACHE → 命中 → 返回 (0.5ms)
降级路径 1:  CACHE → 未命中 → MATERIALIZED → 返回 (3ms)
降级路径 2:  CACHE → 未命中 → MATERIALIZED → 过期 → SQL → 返回 (45ms)
降级路径 3:  全部不可用 → 返回默认值 + 告警 (0ms + alert)

#12.3 背压控制

当系统负载过高时,ComputationCoordinator 会主动限制高成本计算模式:

Python
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 Foundrycoomia-dip
计算引擎Spark (批处理为主)7 种模式混合调度
路由策略固定管道 (Pipeline)动态成本路由
缓存层内置 Object 缓存两级缓存 (Caffeine + Redis)
实时计算Foundry Streaming (受限)Kafka Streams (完整)
物化视图Dataset MaterializationDoris 物化视图
自定义函数TypeScript FunctionsPython 函数计算
多租户路由资源隔离 (Namespace)策略级路由隔离
降级策略手动配置自动熔断 + 降级

coomia-dip 的核心优势在于动态路由——Foundry 的计算路径在 Pipeline 构建时确定,而 coomia-dip 在每次查询时实时选择最优路径。这意味着面对数据增长、负载变化、组件故障等场景,coomia-dip 能自动调整而无需人工干预。

#14. 可观测性与调优

#14.1 路由决策日志

每次路由决策都会记录详细日志,便于分析和优化:

JSON
{
  "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 性能调优建议

  1. 提高缓存命中率:分析 cache miss 日志,对高频属性调整 TTL 或增加预热范围
  2. 物化视图覆盖:对 SQL 模式的高频查询,考虑创建物化视图
  3. 函数计算优化:对 > 200ms 的函数,检查是否可以拆分为缓存 + 增量计算
  4. 减少降级:熔断频繁的引擎需要排查根因(通常是资源不足或网络问题)

#Key Takeaways

  1. 7 种计算模式覆盖了从微秒到秒级的完整延迟光谱。 ComputationCoordinator 不是简单地"尝试缓存,失败就查数据库",而是基于成本模型对 7 种模式进行全局最优选择。每种模式有明确的适用场景、性能特征和保护机制,形成完整的计算策略矩阵。

  2. 动态路由让平台具备自适应能力。 传统系统的计算路径在开发时确定,一旦数据量增长或负载变化就需要人工调整。coomia-dip 的成本模型基于实时指标和历史统计动态调整,路由决策自动适应系统状态变化。结合多租户策略配置,不同业务场景可以获得定制化的计算体验。

  3. 熔断降级保证系统永不因计算引擎故障而完全不可用。 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 #智策平台