返回博客

事件驱动架构:Kafka 在平台中的 7 种角色

TL;DR

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

事件驱动架构:Kafka 在平台中的 7 种角色

系列:S2 架构全景 · 第 6 篇 | 难度:中级 | 阅读时间:18 分钟

TL;DR

  • Kafka 在 coomia-dip 中不只是一个消息队列——它同时扮演 7 种截然不同的角色:CDC 变更捕获、3 主题审计日志、订阅路由、管道触发、派生属性级联、推理触发和跨 Layer 协调。每种角色有独立的 Topic 命名规范、分区策略和消费者组设计。
  • 事件统一采用 Protobuf 编码(而非 JSON),与平台的 gRPC-first 策略一致。单条事件的序列化大小约为 JSON 的 1/3 到 1/5,在高吞吐场景下显著降低网络和存储开销。
  • 通过 Topic 命名约定 {project_id}.{role}.{sub_type} 和 Consumer Group 命名约定 {Layer}-{service}-{role},实现了事件流的可发现性、可追溯性和跨团队零协调。

#1. 引言:消息队列还是事件骨干?

很多平台在技术选型时,把 Kafka 当作"消息队列"——一个异步通信的管道。这种认知局限了 Kafka 的使用场景,通常只用于解耦服务之间的调用。

在 coomia-dip 中,Kafka 的定位完全不同:它是平台的事件骨干(Event Backbone)。每一次数据变更、每一次权限检查、每一次计算请求,都通过 Kafka 事件流串联起来。

Code
传统用法: Service A ──消息──> Kafka ──消息──> Service B
         (异步 RPC 替代品)

coomia-dip 用法:
  ┌───────────┐    ┌─────────────────────────┐
  │ Control Layer   │    │        Kafka            │
  │ Control   │──> │ ┌─────────────────────┐ │
  │           │    │ │ Role 1: CDC         │ │──> Data Layer (Data)
  │           │    │ │ Role 2: Audit (x3)  │ │──> Audit Store
  │           │    │ │ Role 3: Subscription │ │──> External Apps
  │           │    │ │ Role 4: Pipeline     │ │──> Data Layer (Pipeline)
  │           │    │ │ Role 5: Derived Prop │ │──> Reasoning & Decision Layer (Reasoning)
  │           │    │ │ Role 6: Reasoning    │ │──> Reasoning & Decision Layer (Reasoning)
  │           │    │ │ Role 7: Cross-Layer  │ │──> All Layers
  │           │    │ └─────────────────────┘ │
  └───────────┘    └─────────────────────────┘

让我们逐一深入每种角色。

#2. 角色一:CDC 变更数据捕获

#2.1 什么是 CDC

CDC(Change Data Capture)是指捕获数据的每一次创建、更新和删除操作,并将其转化为事件流。在 coomia-dip 中,所有通过 OntologyRuntimeService 的数据变更都会生成 CDC 事件。

#2.2 Topic 设计

Code
Topic 命名: {project_id}.cdc.data-change
分区策略: 按 object_type hash 分区
分区数量: 12 (默认, 可按 Project 规模调整)
保留期:   7 天
压缩策略: delete (超期删除)

#2.3 事件 Schema

PROTOBUF
message DataChangeEvent {
  string event_id = 1;           // UUID, 全局唯一
  string project_id = 2;         // 所属 Project
  string world_id = 3;           // 所属 World
  string object_type = 4;        // 对象类型
  string object_id = 5;          // 对象 ID

  ChangeType change_type = 6;    // CREATE / UPDATE / DELETE
  string principal_id = 7;       // 操作人

  // 变更内容
  map<string, Value> before = 8; // 变更前 (UPDATE/DELETE 时有值)
  map<string, Value> after = 9;  // 变更后 (CREATE/UPDATE 时有值)
  repeated string changed_properties = 10; // 变更的属性列表

  google.protobuf.Timestamp timestamp = 11;
  string correlation_id = 12;    // 请求关联 ID (追踪用)
}

enum ChangeType {
  CREATE = 0;
  UPDATE = 1;
  DELETE = 2;
}

#2.4 CDC 的下游消费者

Code
CDC Topic 的消费者:

Consumer Group 1: data-Layer-materializer
  └── Data Layer: 将变更同步到 Iceberg 表 (数据湖持久化)

Consumer Group 2: data-Layer-search-indexer
  └── Data Layer: 更新 Doris 的倒排索引和向量索引

Consumer Group 3: intelligence-derived-property
  └── Reasoning & Decision Layer: 检测是否需要触发派生属性重计算

Consumer Group 4: control-subscription-router
  └── Control Layer: 检查是否有匹配的订阅规则

Consumer Group 5: intelligence-reasoning-trigger
  └── Reasoning & Decision Layer: 检查是否需要触发规则引擎

一条 CDC 事件会被 5 个不同的 Consumer Group 独立消费——这就是 Kafka 的发布-订阅模型的强大之处。每个消费者按照自己的节奏处理,互不干扰。

#3. 角色二:审计日志(3 个主题)

#3.1 三主题设计

我们不使用单一的审计 Topic——而是根据审计事件的类型和保留需求拆分为 3 个主题:

Code
Topic 1: {project_id}.audit.data-access
  内容: 数据读取操作 (GET, SEARCH, AGGREGATE)
  保留期: 90 天
  分区数: 6
  分区策略: 按 principal_id hash
  典型 QPS: 500-2000 (读操作远多于写操作)

Topic 2: {project_id}.audit.data-mutation
  内容: 数据写入操作 (CREATE, UPDATE, DELETE)
  保留期: 365 天
  分区数: 6
  分区策略: 按 object_type hash
  典型 QPS: 50-200

Topic 3: {project_id}.audit.admin-operation
  内容: 管理操作 (Schema 变更, 权限变更, 配置变更)
  保留期: 永久 (compact 模式)
  分区数: 3
  分区策略: 按 operation_type hash
  典型 QPS: 1-10

#3.2 为什么不用一个 Topic

Code
方案对比:

单一 Topic:
  ├── 优点: 管理简单
  ├── 缺点:
  │   ├── 90 天前的读取审计占用大量存储 (实际已无价值)
  │   ├── 合规审计只需要 mutation 和 admin, 却必须扫描所有事件
  │   ├── 保留期只能设为最长的 (永久), 存储成本爆炸
  │   └── 消费者必须在客户端过滤, 浪费网络带宽
  └── 结论: ❌

三个 Topic:
  ├── 优点:
  │   ├── 每个 Topic 独立保留期, 存储成本最优
  │   ├── 合规审计直接读 mutation + admin Topic, 无需过滤
  │   ├── 实时监控只读 data-access Topic
  │   └── 消费者直接订阅所需 Topic, 零浪费
  ├── 缺点: 管理 3 个 Topic (但 Auto-Provisioning 自动创建)
  └── 结论: ✅

#3.3 审计事件 Schema

PROTOBUF
message AuditEvent {
  string event_id = 1;
  string project_id = 2;
  string world_id = 3;

  // 主体信息
  string principal_id = 4;
  string principal_type = 5;     // USER / SERVICE / SYSTEM
  repeated string principal_roles = 6;
  string source_ip = 7;

  // 操作信息
  string operation = 8;          // GET_OBJECT / SEARCH / CREATE / UPDATE / DELETE / ...
  string object_type = 9;
  string object_id = 10;
  repeated string accessed_properties = 11;

  // 结果
  AuditResult result = 12;       // SUCCESS / DENIED / ERROR
  string denial_reason = 13;     // 如果 DENIED, 说明原因
  int32 result_count = 14;       // 返回结果数量

  // 上下文
  google.protobuf.Timestamp timestamp = 15;
  string correlation_id = 16;
  int64 duration_ms = 17;        // 操作耗时
}

#4. 角色三:订阅路由

#4.1 什么是订阅路由

订阅路由允许外部应用和内部服务订阅特定数据变更。例如:

  • "当设备 A 的温度超过 80°C 时,通知运维系统"
  • "当订单状态变为 SHIPPED 时,触发物流跟踪"
  • "当库存量低于安全线时,通知采购系统"

#4.2 Topic 设计

Code
Topic: {project_id}.subscription.routed
分区数: 6
分区策略: 按 subscription_id hash
保留期: 24 小时 (路由后即消费)

#4.3 订阅规则引擎

Code
订阅规则定义:

{
  "subscription_id": "sub-temp-alert-001",
  "name": "设备温度预警",
  "object_type": "Equipment",
  "filter": {
    "property": "temperature",
    "operator": "GREATER_THAN",
    "value": 80.0
  },
  "trigger_on": ["UPDATE"],
  "notify": {
    "type": "WEBHOOK",
    "url": "https://ops-system.internal/api/alerts",
    "headers": { "Authorization": "Bearer {{token}}" }
  },
  "rate_limit": {
    "max_per_minute": 10,
    "cooldown_seconds": 300
  }
}

#4.4 路由流程

Code
CDC 事件到达 → SubscriptionRouter (Consumer Group: control-subscription-router)
  │
  ├── Step 1: 加载所有匹配的订阅规则
  │   └── 根据 object_type 和 change_type 过滤
  │
  ├── Step 2: 评估每条规则的 filter 条件
  │   └── 检查 changed_properties 是否包含订阅的 property
  │   └── 评估条件表达式 (temperature > 80.0)
  │
  ├── Step 3: 限流检查
  │   └── 检查 rate_limit (该订阅最近 1 分钟内已触发次数)
  │   └── 检查 cooldown (上次触发是否在冷却期内)
  │
  ├── Step 4: 生成路由事件
  │   └── 写入 {project_id}.subscription.routed Topic
  │
  └── Step 5: 执行通知
      └── Webhook / gRPC / 内部事件
      └── 异步执行, 失败重试 (指数退避, 最多 3 次)

#5. 角色四:管道触发

#5.1 场景

数据管道(Pipeline)通常由调度器(DolphinScheduler)定时触发。但有些场景需要事件驱动的管道执行——当上游数据变更时,自动触发下游的 ETL 管道。

#5.2 Topic 设计

Code
Topic: {project_id}.pipeline.trigger
分区数: 6
分区策略: 按 pipeline_id hash
保留期: 48 小时

#5.3 触发规则

PROTOBUF
message PipelineTriggerEvent {
  string event_id = 1;
  string project_id = 2;
  string pipeline_id = 3;        // 要触发的管道
  string trigger_type = 4;       // CDC / SCHEDULE / MANUAL / DEPENDENCY

  // CDC 触发时的信息
  string source_object_type = 5;
  int32 change_count = 6;        // 累积变更数量
  google.protobuf.Timestamp window_start = 7;
  google.protobuf.Timestamp window_end = 8;

  // 触发参数
  map<string, string> parameters = 9;

  google.protobuf.Timestamp timestamp = 10;
}

#5.4 微批次触发策略

为了避免每条 CDC 事件都触发一次管道(太频繁),我们使用微批次策略:

Code
微批次触发策略:

PipelineTriggerAggregator:
  ├── 接收 CDC 事件
  ├── 按 object_type + pipeline_id 分组
  ├── 累积变更直到满足触发条件:
  │   ├── 条件 1: 累积变更 >= threshold (默认 100 条)
  │   ├── 条件 2: 时间窗口 >= interval (默认 5 分钟)
  │   └── 取先满足的条件
  ├── 生成 PipelineTriggerEvent
  └── 写入 pipeline.trigger Topic

示例:
  Pipeline: "equipment-data-etl"
  触发配置: threshold=50, interval=300s

  09:00:00 - 收到 20 条 Equipment 变更 → 累积
  09:02:00 - 收到 15 条 → 累积 = 35
  09:03:30 - 收到 18 条 → 累积 = 53 > 50 → 触发!
  09:03:30 - 发送 PipelineTriggerEvent(change_count=53)
  09:03:31 - Data Layer 的 PipelineExecutor 消费事件, 提交 DolphinScheduler 任务

#6. 角色五:派生属性级联

#6.1 场景

派生属性(Derived Property)的值依赖于其他属性。当依赖的属性变更时,需要触发派生属性的重新计算。如果派生属性 A 依赖于派生属性 B,B 又依赖于属性 C,那么 C 的变更需要级联触发 B 和 A 的重计算。

#6.2 Topic 设计

Code
Topic: {project_id}.derived-property.cascade
分区数: 6
分区策略: 按 object_id hash (确保同一对象的级联事件按序处理)
保留期: 24 小时

#6.3 级联事件 Schema

PROTOBUF
message DerivedPropertyCascadeEvent {
  string event_id = 1;
  string project_id = 2;
  string world_id = 3;
  string object_type = 4;
  string object_id = 5;

  // 触发源
  string trigger_property = 6;       // 变更的源属性
  Value trigger_old_value = 7;
  Value trigger_new_value = 8;

  // 需要重计算的派生属性
  repeated DerivedPropertyTarget targets = 9;

  // 级联层级 (防止无限循环)
  int32 cascade_depth = 10;          // 当前级联深度
  int32 max_cascade_depth = 11;      // 最大允许深度 (默认 10)

  google.protobuf.Timestamp timestamp = 12;
  string correlation_id = 13;
}

message DerivedPropertyTarget {
  string property_name = 1;
  string computation_strategy = 2;   // SQL / EXPRESSION / FUNCTION / ...
  repeated string dependency_properties = 3;
}

#6.4 级联处理流程

Code
CDC 事件: Equipment.temperature 从 75 变为 85

DerivedPropertyCascadeDetector:
  │
  ├── Step 1: 查询依赖图 (DAG)
  │   temperature ← health_score (派生)
  │   health_score ← risk_level (派生)
  │   risk_level ← (无下游依赖)
  │
  ├── Step 2: 拓扑排序
  │   计算顺序: health_score → risk_level
  │
  ├── Step 3: 发送级联事件 (depth=1)
  │   Event 1: 重计算 health_score
  │     trigger: temperature
  │     cascade_depth: 1
  │
  └── Step 4: health_score 计算完成后
      │
      └── 发送级联事件 (depth=2)
          Event 2: 重计算 risk_level
            trigger: health_score
            cascade_depth: 2

安全机制:
  - cascade_depth > max_cascade_depth → 停止级联, 记录告警
  - 循环依赖检测: 在 DAG 构建时检测, 拒绝创建循环依赖的派生属性
  - 每个级联事件包含 correlation_id, 可追踪完整的级联链路

#7. 角色六:推理触发

#7.1 场景

Reasoning & Decision Layer 的推理引擎(规则引擎 + AI 推理)需要在特定条件下自动触发。例如:

  • 当设备的 risk_level 变为 HIGH 时,触发故障诊断推理
  • 当库存低于安全线时,触发补货建议推理
  • 当金融交易金额异常时,触发反欺诈规则链

#7.2 Topic 设计

Code
Topic: {project_id}.reasoning.trigger
分区数: 6
分区策略: 按 reasoning_type hash
保留期: 48 小时

#7.3 推理触发事件

PROTOBUF
message ReasoningTriggerEvent {
  string event_id = 1;
  string project_id = 2;
  string world_id = 3;

  // 触发对象
  string object_type = 4;
  string object_id = 5;

  // 推理配置
  string reasoning_type = 6;         // RULE_CHAIN / AI_INFERENCE / HYBRID
  string reasoning_config_id = 7;    // 推理配置 ID

  // 触发上下文
  map<string, Value> trigger_context = 8;  // 传递给推理引擎的上下文
  string trigger_source = 9;         // CDC / SCHEDULE / MANUAL / CASCADE

  // 优先级
  Priority priority = 10;           // LOW / NORMAL / HIGH / CRITICAL

  google.protobuf.Timestamp timestamp = 11;
  string correlation_id = 12;
}

#7.4 推理触发与 CDC 的关系

Code
CDC 事件 → DerivedPropertyCascade → 派生属性更新
                                         │
                                         ▼
                                   ReasoningTrigger
                                   (基于更新后的派生属性值)

示例:
  1. temperature 从 75 变为 85           (CDC)
  2. health_score 从 0.9 重计算为 0.6    (派生属性级联)
  3. risk_level 从 LOW 重计算为 HIGH     (派生属性级联)
  4. 触发故障诊断推理                     (推理触发)
     → 规则引擎评估 15 条故障规则
     → 输出: "疑似轴承过热, 建议检查润滑系统"
     → 创建 Action: MaintenanceWorkOrder

#8. 角色七:跨 Layer 协调

#8.1 场景

coomia-dip 有 8 个 Layer(实际 5 个独立部署单元)。它们之间的主要通信使用 gRPC(请求-响应),但某些场景需要异步协调:

  • Schema 变更需要通知所有 Layer 刷新缓存
  • 新 World 创建需要多个 Layer 初始化资源
  • 系统维护需要通知所有 Layer 进入只读模式

#8.2 Topic 设计

Code
Topic: platform.coordination.{event_type}
分区数: 3 (全局 Topic, 不按 Project 分区)
保留期: 24 小时

#8.3 协调事件类型

Code
事件类型:

1. schema.changed
   触发: Schema 发生变更 (新增/修改/删除 ObjectType)
   消费者: 所有 Layer 的 Schema 缓存
   效果: 失效本地缓存, 重新加载

2. world.created / world.deleted
   触发: World 生命周期变更
   消费者: Data Layer (Doris 连接池), Reasoning & Decision Layer (计算上下文)
   效果: 初始化/清理 World 相关资源

3. system.maintenance.enter / system.maintenance.exit
   触发: 系统维护窗口
   消费者: 所有 Layer
   效果: 进入/退出只读模式

4. config.updated
   触发: 平台配置变更 (限流阈值, 功能开关)
   消费者: 所有 Layer
   效果: 热更新配置

5. permission.policy.changed
   触发: RBAC/ABAC 策略变更
   消费者: Control Layer (PolicyEngine), 所有带缓存的服务
   效果: 重编译权限决策树

#8.4 协调事件的幂等性

跨 Layer 协调事件必须是幂等的——即使同一事件被消费多次,效果也只发生一次:

Code
幂等性保证:

每个协调事件包含:
  - event_id (UUID): 全局唯一标识
  - event_version (int64): 单调递增版本号

消费者端:
  - 维护 processed_events 集合 (Redis SET, TTL = 24h)
  - 消费前检查: if event_id in processed_events → skip
  - 消费后记录: processed_events.add(event_id)

版本检查:
  - 维护 latest_version (Redis KV)
  - 如果 event_version <= latest_version → skip (旧事件)
  - 如果 event_version > latest_version + 1 → 缺少中间事件, 触发全量同步

#9. Topic 命名规范

#9.1 完整命名约定

Code
命名模式: {scope}.{role}.{sub_type}

scope:
  - {project_id}: Project 级事件 (大多数事件)
  - platform: 平台级全局事件

role:
  - cdc: 变更数据捕获
  - audit: 审计日志
  - subscription: 订阅路由
  - pipeline: 管道触发
  - derived-property: 派生属性级联
  - reasoning: 推理触发
  - coordination: 跨 Layer 协调

sub_type:
  - 具体的事件子类型

完整示例:
  proj_001.cdc.data-change
  proj_001.audit.data-access
  proj_001.audit.data-mutation
  proj_001.audit.admin-operation
  proj_001.subscription.routed
  proj_001.pipeline.trigger
  proj_001.derived-property.cascade
  proj_001.reasoning.trigger
  platform.coordination.schema.changed
  platform.coordination.world.created
  platform.coordination.system.maintenance.enter

#9.2 Consumer Group 命名约定

Code
命名模式: {Layer}-{service}-{role}

示例:
  control-ontology-runtime-cdc          (Control Layer 消费 CDC)
  control-subscription-router-cdc       (Control Layer 订阅路由器)
  data-materializer-cdc                 (Data Layer 物化消费 CDC)
  data-search-indexer-cdc               (Data Layer 搜索索引消费 CDC)
  data-pipeline-executor-trigger        (Data Layer 管道执行器)
  intelligence-derived-property-cascade  (Reasoning & Decision Layer 派生属性级联)
  intelligence-reasoning-engine-trigger  (Reasoning & Decision Layer 推理引擎)
  intelligence-rule-engine-trigger       (Reasoning & Decision Layer 规则引擎)

#10. 事件 Schema 管理

#10.1 为什么用 Protobuf 而不是 JSON

Code
JSON vs Protobuf 对比:

序列化大小:
  JSON DataChangeEvent:  ~800 bytes
  Protobuf DataChangeEvent: ~200 bytes
  压缩比: 4:1

序列化/反序列化性能:
  JSON (Jackson):    ~5μs / ~8μs
  Protobuf:          ~1μs / ~1.5μs
  提升: 5x

Schema 演进:
  JSON: 无强制 Schema, 向后兼容靠约定
  Protobuf: 强类型 Schema, 编译时检查, 向后兼容有保证
    - 新增字段: ✅ (旧消费者忽略)
    - 删除字段: ✅ (标记 reserved)
    - 修改类型: ❌ (编译错误)

跨语言:
  JSON: 每种语言需要手写或生成 DTO
  Protobuf: 一份 .proto 文件生成 Java + Python + TypeScript

#10.2 Protobuf Schema 管理策略

Code
目录结构:

proto/
├── events/
│   ├── cdc.proto                    # CDC 事件
│   ├── audit.proto                  # 审计事件
│   ├── subscription.proto           # 订阅路由事件
│   ├── pipeline.proto               # 管道触发事件
│   ├── derived_property.proto       # 派生属性级联事件
│   ├── reasoning.proto              # 推理触发事件
│   └── coordination.proto           # 跨 Layer 协调事件
├── common/
│   ├── value.proto                  # 通用值类型
│   └── context.proto                # WorldContext 定义
└── buf.yaml                         # Buf Schema Registry 配置

版本管理:
  - 每个 .proto 文件有 package 版本: onto.events.v1
  - 向后兼容的变更: 在同一 package 内添加字段
  - 不兼容的变更: 创建新 package: onto.events.v2
  - 消费者同时支持 v1 和 v2, 通过事件头部的版本标识路由

#10.3 事件信封(Event Envelope)

所有 Kafka 事件都使用统一的信封格式:

PROTOBUF
message EventEnvelope {
  // 信封元数据 (Kafka Headers 中)
  string event_type = 1;       // "cdc.data-change" / "audit.data-access" / ...
  string schema_version = 2;   // "v1" / "v2"
  string content_type = 3;     // "application/x-protobuf"
  string correlation_id = 4;   // 请求追踪 ID
  string source_plane = 5;     // "control" / "data" / "intelligence"

  // 事件体 (Kafka Value 中, Protobuf 编码)
  bytes payload = 6;           // 实际事件内容 (根据 event_type 反序列化)
}

信封的 event_typeschema_version 存储在 Kafka Headers 中,消费者可以在不反序列化 payload 的情况下决定路由策略。

#11. 故障处理与可靠性

#11.1 Dead Letter Topic

当事件处理失败超过重试次数后,事件被路由到 Dead Letter Topic:

Code
Dead Letter Topic: {original_topic}.dlq

处理流程:
  Event → Consumer → 处理失败
    → 重试 1 (1s 后)
    → 重试 2 (5s 后)
    → 重试 3 (30s 后)
    → 路由到 DLQ Topic
    → 告警通知

DLQ 事件格式:
  原始事件 + 错误信息 + 重试历史 + 原始 Topic + 原始 Partition + Offset

DLQ 处理:
  1. 自动: 定时任务每小时检查 DLQ, 尝试重新处理
  2. 手动: 运维人员查看 DLQ, 修复问题后重放事件

#11.2 背压处理

Code
背压策略:

当消费者处理速度 < 生产者发送速度时:

Level 1: Consumer Lag 监控
  if consumer_lag > threshold_warning (1000 条):
    → 告警通知
    → 增加 Consumer 实例数量

Level 2: 动态分区再平衡
  if consumer_lag > threshold_critical (10000 条):
    → 触发 Consumer Group 再平衡
    → 将更多 Partition 分配给空闲 Consumer

Level 3: 生产者限流
  if consumer_lag > threshold_emergency (100000 条):
    → OntologyRuntimeService 降低写入速率
    → 返回 429 Too Many Requests
    → 等待消费者追上

#11.3 事件顺序保证

Code
顺序保证级别:

Level 1: 同一 Partition 内严格有序
  保证: ✅ Kafka 原生保证
  应用: 同一对象的 CDC 事件在同一 Partition (按 object_id hash)

Level 2: 因果有序
  保证: ✅ 通过 correlation_id 和 cascade_depth
  应用: 派生属性级联按 depth 顺序处理

Level 3: 全局有序
  保证: ❌ 不保证, 也不需要
  原因: 不同对象的变更本身就是并行的, 强制全局有序会严重限制吞吐量

#12. 性能基准与调优

#12.1 吞吐量基准

Code
测试环境: 3 Broker, 每 Broker 4 CPU / 16GB RAM

CDC Topic (12 partitions):
  生产者吞吐: 50,000 events/s
  消费者吞吐: 30,000 events/s (单 Consumer Group)
  端到端延迟 P99: 12ms

Audit Topic (6 partitions):
  生产者吞吐: 10,000 events/s
  消费者吞吐: 20,000 events/s (审计存储写入很快)
  端到端延迟 P99: 5ms

Derived Property Cascade (6 partitions):
  生产者吞吐: 5,000 events/s
  消费者吞吐: 2,000 events/s (计算是瓶颈)
  端到端延迟 P99: 50ms (包含计算时间)

#12.2 关键调优参数

Code
Producer 配置:
  acks=1                    # 审计用 acks=all
  batch.size=32768           # 32KB 批量
  linger.ms=5                # 5ms 等待批量
  compression.type=snappy    # Snappy 压缩
  max.in.flight.per.connection=5

Consumer 配置:
  max.poll.records=500       # 每次 poll 最多 500 条
  max.poll.interval.ms=300000 # 5 分钟超时
  session.timeout.ms=30000   # 30 秒会话超时
  auto.offset.reset=latest   # CDC 用 latest, Audit 用 earliest
  enable.auto.commit=false   # 手动提交 offset

#13. 与 Palantir Foundry 的对比

维度Palantir Foundrycoomia-dip
事件系统Foundry 内部事件总线 (闭源)Kafka (开源)
事件格式JSON / AvroProtobuf
CDCDataset Transaction LogCDC Topic per Project
审计Audit Service (集中式)3 Topic 分级审计
派生属性TypeScript OSDK 触发Kafka 级联事件
订阅Webhook + Object Set Subscription规则引擎 + 路由 Topic
跨服务协调内部 RPCKafka 协调 Topic

#Key Takeaways

  1. Kafka 不只是消息队列——它是平台的事件骨干。 通过 7 种角色的划分,coomia-dip 将 Kafka 从一个简单的异步通信管道提升为平台级的事件编排系统。每种角色有独立的 Topic、分区策略和消费者组,互不干扰,独立扩缩。

  2. Protobuf 事件 Schema 是跨语言协作的基础。 Java (Control Layer + Data Layer) 和 Python (Reasoning & Decision Layer + Agent Runtime Layer) 共享同一套 .proto 定义,编译时类型检查确保事件生产者和消费者始终兼容。相比 JSON,Protobuf 在序列化大小和性能上有 4-5 倍的优势,在高吞吐场景下尤为关键。

  3. 级联事件的深度控制是安全阀。 派生属性级联和推理触发都可能形成长链路(A → B → C → ...)。cascade_depth 字段和 max_cascade_depth 限制确保级联不会失控。循环依赖在 DAG 构建阶段就被检测和拒绝,而不是在运行时发现无限循环。

下一篇预告: [S2-07] 计算策略路由:7 种计算模式的优先级调度——深入理解 ComputationCoordinator 如何在 SQL、表达式、Reducer、函数、缓存、物化和实时计算之间做出最优选择。

Tags: #event-driven #kafka #cdc #audit #subscription #pipeline #derived-property #reasoning #protobuf #coomia-dip #智策平台