事件驱动架构:Kafka 在平台中的 7 种角色
TL;DR
事件驱动架构: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 事件流串联起来。
传统用法: 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 设计
Topic 命名: {project_id}.cdc.data-change
分区策略: 按 object_type hash 分区
分区数量: 12 (默认, 可按 Project 规模调整)
保留期: 7 天
压缩策略: delete (超期删除)
#2.3 事件 Schema
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 的下游消费者
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 个主题:
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
方案对比:
单一 Topic:
├── 优点: 管理简单
├── 缺点:
│ ├── 90 天前的读取审计占用大量存储 (实际已无价值)
│ ├── 合规审计只需要 mutation 和 admin, 却必须扫描所有事件
│ ├── 保留期只能设为最长的 (永久), 存储成本爆炸
│ └── 消费者必须在客户端过滤, 浪费网络带宽
└── 结论: ❌
三个 Topic:
├── 优点:
│ ├── 每个 Topic 独立保留期, 存储成本最优
│ ├── 合规审计直接读 mutation + admin Topic, 无需过滤
│ ├── 实时监控只读 data-access Topic
│ └── 消费者直接订阅所需 Topic, 零浪费
├── 缺点: 管理 3 个 Topic (但 Auto-Provisioning 自动创建)
└── 结论: ✅
#3.3 审计事件 Schema
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 设计
Topic: {project_id}.subscription.routed
分区数: 6
分区策略: 按 subscription_id hash
保留期: 24 小时 (路由后即消费)
#4.3 订阅规则引擎
订阅规则定义:
{
"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 路由流程
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 设计
Topic: {project_id}.pipeline.trigger
分区数: 6
分区策略: 按 pipeline_id hash
保留期: 48 小时
#5.3 触发规则
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 事件都触发一次管道(太频繁),我们使用微批次策略:
微批次触发策略:
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 设计
Topic: {project_id}.derived-property.cascade
分区数: 6
分区策略: 按 object_id hash (确保同一对象的级联事件按序处理)
保留期: 24 小时
#6.3 级联事件 Schema
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 级联处理流程
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 设计
Topic: {project_id}.reasoning.trigger
分区数: 6
分区策略: 按 reasoning_type hash
保留期: 48 小时
#7.3 推理触发事件
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 的关系
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 设计
Topic: platform.coordination.{event_type}
分区数: 3 (全局 Topic, 不按 Project 分区)
保留期: 24 小时
#8.3 协调事件类型
事件类型:
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 协调事件必须是幂等的——即使同一事件被消费多次,效果也只发生一次:
幂等性保证:
每个协调事件包含:
- 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 完整命名约定
命名模式: {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 命名约定
命名模式: {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
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 管理策略
目录结构:
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 事件都使用统一的信封格式:
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_type 和 schema_version 存储在 Kafka Headers 中,消费者可以在不反序列化 payload 的情况下决定路由策略。
#11. 故障处理与可靠性
#11.1 Dead Letter Topic
当事件处理失败超过重试次数后,事件被路由到 Dead Letter Topic:
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 背压处理
背压策略:
当消费者处理速度 < 生产者发送速度时:
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 事件顺序保证
顺序保证级别:
Level 1: 同一 Partition 内严格有序
保证: ✅ Kafka 原生保证
应用: 同一对象的 CDC 事件在同一 Partition (按 object_id hash)
Level 2: 因果有序
保证: ✅ 通过 correlation_id 和 cascade_depth
应用: 派生属性级联按 depth 顺序处理
Level 3: 全局有序
保证: ❌ 不保证, 也不需要
原因: 不同对象的变更本身就是并行的, 强制全局有序会严重限制吞吐量
#12. 性能基准与调优
#12.1 吞吐量基准
测试环境: 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 关键调优参数
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 Foundry | coomia-dip |
|---|---|---|
| 事件系统 | Foundry 内部事件总线 (闭源) | Kafka (开源) |
| 事件格式 | JSON / Avro | Protobuf |
| CDC | Dataset Transaction Log | CDC Topic per Project |
| 审计 | Audit Service (集中式) | 3 Topic 分级审计 |
| 派生属性 | TypeScript OSDK 触发 | Kafka 级联事件 |
| 订阅 | Webhook + Object Set Subscription | 规则引擎 + 路由 Topic |
| 跨服务协调 | 内部 RPC | Kafka 协调 Topic |
#Key Takeaways
-
Kafka 不只是消息队列——它是平台的事件骨干。 通过 7 种角色的划分,coomia-dip 将 Kafka 从一个简单的异步通信管道提升为平台级的事件编排系统。每种角色有独立的 Topic、分区策略和消费者组,互不干扰,独立扩缩。
-
Protobuf 事件 Schema 是跨语言协作的基础。 Java (Control Layer + Data Layer) 和 Python (Reasoning & Decision Layer + Agent Runtime Layer) 共享同一套 .proto 定义,编译时类型检查确保事件生产者和消费者始终兼容。相比 JSON,Protobuf 在序列化大小和性能上有 4-5 倍的优势,在高吞吐场景下尤为关键。
-
级联事件的深度控制是安全阀。 派生属性级联和推理触发都可能形成长链路(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 #智策平台