事件订阅与通知:让平台主动告诉你
在前面的教程中,我们掌握了如何主动查询数据、执行 Action、触发规则。但在真实业务场景中,"被动等通知"往往比"主动去轮询"更高效——当订单状态变化时自动通知下游系统,当指标异常时即时推送告警,当审批完成时触发后续流程。
“系列:S12 开发者教程 · 第 13 篇 | 难度:中级 | 阅读时间:15 分钟
事件订阅与通知:让平台主动告诉你
#引言
在前面的教程中,我们掌握了如何主动查询数据、执行 Action、触发规则。但在真实业务场景中,"被动等通知"往往比"主动去轮询"更高效——当订单状态变化时自动通知下游系统,当指标异常时即时推送告警,当审批完成时触发后续流程。
coomia-dip 的事件订阅系统正是为此设计的。它基于本体模型中的对象变更事件,提供了声明式的订阅机制。你只需要告诉平台"我关心什么",平台就会在事件发生时主动将消息推送给你。本教程将带你从零搭建一套完整的事件订阅与通知体系。
#1. 理解 coomia-dip 事件模型
#1.1 事件的本质
在 coomia-dip 中,一切变更都是事件。当你通过 Action 修改一个对象的属性、创建一个新的关联关系、或者删除一个对象时,平台内部都会产生对应的变更事件。这些事件构成了一条不间断的事件流(Event Stream),订阅系统的工作就是从这条流中筛选出你关心的事件并投递给你。
from ontology_sdk.models import OntologyEvent
# 事件的基本结构
class OntologyEvent:
"""本体变更事件"""
event_id: str # 全局唯一事件ID
event_type: EventType # CREATE / UPDATE / DELETE / LINK / UNLINK
object_type: str # 发生变更的对象类型
object_rid: str # 发生变更的对象 RID
timestamp: datetime # 事件发生时间戳(UTC)
changes: dict[str, Change] # 字段级变更详情
actor: str # 触发变更的用户或系统
transaction_id: str # 所属事务ID
#1.2 事件类型详解
coomia-dip 将事件分为五种基本类型:
| 事件类型 | 含义 | 触发场景 |
|---|---|---|
CREATE | 新对象创建 | 通过 Action 创建对象、数据导入 |
UPDATE | 对象属性变更 | 修改对象的一个或多个属性 |
DELETE | 对象删除 | 软删除或硬删除 |
LINK | 建立关联关系 | 两个对象之间创建 Link |
UNLINK | 解除关联关系 | 两个对象之间移除 Link |
每种事件类型都携带完整的变更上下文,包括变更前的值(before)和变更后的值(after),让订阅者可以精确感知发生了什么。
#1.3 事件传播的保证
coomia-dip 的事件系统提供以下保证:
- 至少一次投递(At-Least-Once):每个事件至少会被投递一次到每个活跃订阅
- 事件有序性:同一个对象的事件保证按照发生顺序投递
- 持久化存储:事件在投递前会持久化到事件存储,即使系统重启也不会丢失
- 重放能力:支持从历史某个时间点开始重放事件
# 事件存储的持久化机制
class EventStore:
"""
事件存储采用 Apache Kafka 作为底层实现:
- 每个 ObjectType 对应一个 Kafka Topic
- 分区键为 object_rid,保证同一对象的事件有序
- 保留策略默认 7 天,可按需配置
"""
def append(self, event: OntologyEvent) -> None:
topic = f"onto.events.{event.object_type}"
partition_key = event.object_rid
self.kafka_producer.send(topic, key=partition_key, value=event.serialize())
def replay(self, object_type: str, from_offset: int) -> Iterator[OntologyEvent]:
"""从指定 offset 开始重放事件"""
consumer = self.kafka_consumer.subscribe(f"onto.events.{object_type}")
consumer.seek(from_offset)
for record in consumer:
yield OntologyEvent.deserialize(record.value)
#2. 创建你的第一个订阅
#2.1 通过 SDK 创建订阅
让我们从一个实际场景开始:你需要在"工单"对象的状态变为"已完成"时收到通知。
from ontology_sdk import OntoPlatform
from ontology_sdk.subscription import (
Subscription,
EventFilter,
FilterCondition,
WebhookTarget,
)
# 初始化平台连接
platform = OntoPlatform(
base_url="http://localhost:8080",
token="your-api-token",
)
# 定义订阅
subscription = Subscription(
name="ticket-completed-notify",
description="工单完成时通知下游系统",
object_type="Ticket",
event_types=["UPDATE"],
filter=EventFilter(
conditions=[
FilterCondition(
field="status",
operator="eq",
value="COMPLETED",
apply_to="after", # 对变更后的值进行过滤
),
FilterCondition(
field="status",
operator="neq",
value="COMPLETED",
apply_to="before", # 变更前不是 COMPLETED(排除重复触发)
),
],
logic="AND",
),
target=WebhookTarget(
url="https://your-service.example.com/webhook/ticket-completed",
method="POST",
headers={"Authorization": "Bearer your-webhook-secret"},
retry_policy={"max_retries": 3, "backoff_ms": 1000},
),
)
# 注册订阅
result = platform.subscriptions.create(subscription)
print(f"订阅创建成功: {result.subscription_id}")
print(f"状态: {result.status}")
#2.2 订阅的过滤条件
过滤条件是订阅系统的核心能力。coomia-dip 支持丰富的过滤表达式:
# 示例 1:多条件组合过滤
filter_complex = EventFilter(
conditions=[
FilterCondition(field="priority", operator="in", value=["HIGH", "CRITICAL"]),
FilterCondition(field="assignee", operator="is_not_null"),
FilterCondition(field="updated_at", operator="gt", value="2025-01-01T00:00:00Z"),
],
logic="AND",
)
# 示例 2:嵌套过滤(OR + AND)
filter_nested = EventFilter(
logic="OR",
groups=[
EventFilter(
logic="AND",
conditions=[
FilterCondition(field="status", operator="eq", value="CRITICAL", apply_to="after"),
FilterCondition(field="region", operator="eq", value="APAC"),
],
),
EventFilter(
logic="AND",
conditions=[
FilterCondition(field="status", operator="eq", value="DOWN", apply_to="after"),
FilterCondition(field="sla_tier", operator="eq", value="PLATINUM"),
],
),
],
)
# 示例 3:变更幅度过滤
filter_delta = EventFilter(
conditions=[
FilterCondition(
field="temperature",
operator="delta_gt", # 变化幅度大于
value=5.0,
apply_to="delta", # 比较 |after - before|
),
],
)
#2.3 投递目标类型
coomia-dip 支持多种投递目标:
from ontology_sdk.subscription import (
WebhookTarget,
GrpcTarget,
KafkaTarget,
InternalActionTarget,
)
# 1. Webhook — 最通用的方式
webhook = WebhookTarget(
url="https://api.example.com/events",
method="POST",
headers={"Content-Type": "application/json"},
timeout_ms=5000,
retry_policy={"max_retries": 3, "backoff_ms": 1000},
)
# 2. gRPC — 高性能场景
grpc_target = GrpcTarget(
endpoint="event-handler.internal:50051",
service="EventHandlerService",
method="HandleOntologyEvent",
tls_enabled=True,
)
# 3. Kafka — 大规模异步处理
kafka_target = KafkaTarget(
bootstrap_servers="kafka:9092",
topic="downstream.events",
key_expression="${event.object_rid}",
partition_strategy="BY_KEY",
)
# 4. 内部 Action — 链式触发
action_target = InternalActionTarget(
action_type="NotifyStakeholders",
parameter_mapping={
"ticket_id": "${event.object_rid}",
"new_status": "${event.changes.status.after}",
"changed_by": "${event.actor}",
},
)
#3. gRPC 流式订阅
#3.1 服务端推流(Server Streaming)
对于需要低延迟、持续接收事件的场景,coomia-dip 提供了 gRPC 流式订阅接口:
// subscription.proto
syntax = "proto3";
package onto.subscription.v1;
service SubscriptionService {
// 创建订阅
rpc CreateSubscription(CreateSubscriptionRequest) returns (CreateSubscriptionResponse);
// 流式接收事件
rpc StreamEvents(StreamEventsRequest) returns (stream OntologyEventMessage);
// 确认事件已处理
rpc AcknowledgeEvents(AcknowledgeRequest) returns (AcknowledgeResponse);
// 管理订阅
rpc ListSubscriptions(ListRequest) returns (ListResponse);
rpc PauseSubscription(PauseRequest) returns (PauseResponse);
rpc ResumeSubscription(ResumeRequest) returns (ResumeResponse);
rpc DeleteSubscription(DeleteRequest) returns (DeleteResponse);
}
message StreamEventsRequest {
string subscription_id = 1;
int64 from_offset = 2; // 可选,从指定位置开始
int32 batch_size = 3; // 每批事件数量
int32 heartbeat_interval_ms = 4; // 心跳间隔
}
message OntologyEventMessage {
string event_id = 1;
string event_type = 2;
string object_type = 3;
string object_rid = 4;
int64 timestamp_ms = 5;
map<string, FieldChange> changes = 6;
string actor = 7;
string transaction_id = 8;
int64 offset = 9; // 用于 ACK
}
#3.2 Python 客户端实现
import grpc
import asyncio
from ontology_sdk.grpc_client import SubscriptionServiceStub
async def stream_events():
"""流式接收工单变更事件"""
channel = grpc.aio.insecure_channel("localhost:50051")
stub = SubscriptionServiceStub(channel)
request = StreamEventsRequest(
subscription_id="sub-ticket-completed",
batch_size=10,
heartbeat_interval_ms=30000,
)
ack_batch = []
async for event in stub.StreamEvents(request):
# 心跳消息(空事件)
if not event.event_id:
print("Heartbeat received")
continue
# 处理事件
print(f"[{event.event_type}] {event.object_type}/{event.object_rid}")
print(f" Changes: {dict(event.changes)}")
print(f" Actor: {event.actor}")
# 批量确认
ack_batch.append(event.offset)
if len(ack_batch) >= 10:
await stub.AcknowledgeEvents(AcknowledgeRequest(
subscription_id="sub-ticket-completed",
offsets=ack_batch,
))
ack_batch.clear()
print(" Acknowledged batch")
await channel.close()
# 运行
asyncio.run(stream_events())
#3.3 带重连的健壮客户端
生产环境中,网络断开和服务重启是常态。以下是一个带自动重连的健壮实现:
import grpc
import asyncio
import logging
from typing import Callable, Awaitable
logger = logging.getLogger(__name__)
class ResilientEventSubscriber:
"""带自动重连的事件订阅客户端"""
def __init__(
self,
endpoint: str,
subscription_id: str,
handler: Callable[[OntologyEventMessage], Awaitable[None]],
max_retries: int = -1, # -1 表示无限重试
base_backoff_s: float = 1.0,
max_backoff_s: float = 60.0,
):
self.endpoint = endpoint
self.subscription_id = subscription_id
self.handler = handler
self.max_retries = max_retries
self.base_backoff_s = base_backoff_s
self.max_backoff_s = max_backoff_s
self._last_offset: int = 0
self._running = False
async def start(self):
"""启动订阅,自动重连"""
self._running = True
retries = 0
while self._running:
try:
channel = grpc.aio.insecure_channel(self.endpoint)
stub = SubscriptionServiceStub(channel)
request = StreamEventsRequest(
subscription_id=self.subscription_id,
from_offset=self._last_offset,
batch_size=20,
heartbeat_interval_ms=30000,
)
logger.info(f"Connecting to event stream from offset {self._last_offset}")
async for event in stub.StreamEvents(request):
retries = 0 # 成功收到消息,重置重试计数
if not event.event_id:
continue # 心跳
await self.handler(event)
self._last_offset = event.offset + 1
# 逐条确认(可优化为批量)
await stub.AcknowledgeEvents(AcknowledgeRequest(
subscription_id=self.subscription_id,
offsets=[event.offset],
))
except grpc.aio.AioRpcError as e:
if e.code() == grpc.StatusCode.UNAVAILABLE:
logger.warning(f"Server unavailable, will retry: {e.details()}")
elif e.code() == grpc.StatusCode.NOT_FOUND:
logger.error(f"Subscription not found: {self.subscription_id}")
break
else:
logger.error(f"gRPC error: {e.code()} - {e.details()}")
except Exception as e:
logger.error(f"Unexpected error: {e}")
finally:
try:
await channel.close()
except Exception:
pass
if not self._running:
break
retries += 1
if self.max_retries >= 0 and retries > self.max_retries:
logger.error("Max retries exceeded, stopping subscriber")
break
backoff = min(self.base_backoff_s * (2 ** (retries - 1)), self.max_backoff_s)
logger.info(f"Reconnecting in {backoff:.1f}s (attempt {retries})")
await asyncio.sleep(backoff)
async def stop(self):
"""优雅停止"""
self._running = False
# 使用示例
async def handle_event(event: OntologyEventMessage):
print(f"Processing: {event.event_type} on {event.object_rid}")
# 你的业务逻辑...
subscriber = ResilientEventSubscriber(
endpoint="localhost:50051",
subscription_id="sub-ticket-completed",
handler=handle_event,
)
# 启动
await subscriber.start()
#4. 实战场景:多系统联动
#4.1 场景描述
假设你正在构建一个 IT 服务管理系统(ITSM),需要实现以下联动:
- 工单创建 → 通知 Slack 频道
- 工单分配 → 发送邮件给负责人
- 工单完成 → 触发满意度调查
- SLA 即将到期 → 升级告警
#4.2 实现方案
from ontology_sdk import OntoPlatform
from ontology_sdk.subscription import Subscription, EventFilter, FilterCondition
platform = OntoPlatform(base_url="http://localhost:8080", token="your-token")
# 订阅 1:工单创建通知
sub_created = Subscription(
name="ticket-created-slack",
object_type="Ticket",
event_types=["CREATE"],
target=WebhookTarget(
url="https://hooks.slack.com/services/YOUR/SLACK/WEBHOOK",
method="POST",
body_template="""{
"channel": "#itsm-tickets",
"text": "新工单 <${event.object_rid}|${event.changes.title.after}> 已创建\\n优先级: ${event.changes.priority.after}\\n创建人: ${event.actor}"
}""",
),
)
# 订阅 2:工单分配通知
sub_assigned = Subscription(
name="ticket-assigned-email",
object_type="Ticket",
event_types=["UPDATE"],
filter=EventFilter(
conditions=[
FilterCondition(field="assignee", operator="is_not_null", apply_to="after"),
FilterCondition(field="assignee", operator="changed"),
],
logic="AND",
),
target=WebhookTarget(
url="https://email-service.internal/send",
method="POST",
body_template="""{
"to": "${event.changes.assignee.after}",
"subject": "新工单分配: ${event.changes.title.after}",
"body": "你被分配了工单 ${event.object_rid},优先级 ${event.changes.priority.after}"
}""",
),
)
# 订阅 3:工单完成触发调查
sub_completed = Subscription(
name="ticket-completed-survey",
object_type="Ticket",
event_types=["UPDATE"],
filter=EventFilter(
conditions=[
FilterCondition(field="status", operator="eq", value="COMPLETED", apply_to="after"),
FilterCondition(field="status", operator="neq", value="COMPLETED", apply_to="before"),
],
logic="AND",
),
target=InternalActionTarget(
action_type="SendSatisfactionSurvey",
parameter_mapping={
"ticket_id": "${event.object_rid}",
"reporter": "${event.changes.reporter.after}",
"resolved_by": "${event.changes.assignee.after}",
},
),
)
# 订阅 4:SLA 告警(通过定时检查 + 事件触发的混合方式)
sub_sla_warning = Subscription(
name="ticket-sla-warning",
object_type="Ticket",
event_types=["UPDATE"],
filter=EventFilter(
conditions=[
FilterCondition(field="sla_remaining_hours", operator="lt", value=2),
FilterCondition(field="status", operator="in", value=["OPEN", "IN_PROGRESS"]),
],
logic="AND",
),
target=WebhookTarget(
url="https://pagerduty.example.com/v2/enqueue",
method="POST",
body_template="""{
"routing_key": "YOUR_PD_KEY",
"event_action": "trigger",
"payload": {
"summary": "SLA Warning: Ticket ${event.object_rid} has <2h remaining",
"severity": "warning",
"source": "coomia-dip"
}
}""",
),
)
# 批量注册
for sub in [sub_created, sub_assigned, sub_completed, sub_sla_warning]:
result = platform.subscriptions.create(sub)
print(f"Created: {result.subscription_id} -> {result.status}")
#5. 订阅管理与监控
#5.1 查看订阅状态
# 列出所有订阅
subscriptions = platform.subscriptions.list()
for sub in subscriptions:
print(f"{sub.name}: {sub.status}")
print(f" Created: {sub.created_at}")
print(f" Events delivered: {sub.stats.total_delivered}")
print(f" Events failed: {sub.stats.total_failed}")
print(f" Last event at: {sub.stats.last_event_at}")
print(f" Lag: {sub.stats.current_lag} events")
print()
# 查看特定订阅的详细统计
stats = platform.subscriptions.get_stats("sub-ticket-completed")
print(f"投递成功率: {stats.success_rate:.1%}")
print(f"平均延迟: {stats.avg_latency_ms:.0f}ms")
print(f"P99 延迟: {stats.p99_latency_ms:.0f}ms")
#5.2 暂停与恢复
# 暂停订阅(维护时使用)
platform.subscriptions.pause("sub-ticket-completed")
print("Subscription paused")
# 恢复订阅(会从暂停点继续投递)
platform.subscriptions.resume("sub-ticket-completed")
print("Subscription resumed, catching up...")
# 重置订阅偏移量(重新处理历史事件)
platform.subscriptions.reset_offset(
"sub-ticket-completed",
to_timestamp="2025-01-01T00:00:00Z",
)
print("Offset reset, replaying from 2025-01-01")
#5.3 死信队列(DLQ)
投递失败超过重试次数的事件不会丢失,而是进入死信队列:
# 查看死信队列中的事件
dead_letters = platform.subscriptions.list_dead_letters("sub-ticket-completed")
for dl in dead_letters:
print(f"Event: {dl.event_id}")
print(f" Failed at: {dl.failed_at}")
print(f" Attempts: {dl.attempt_count}")
print(f" Last error: {dl.last_error}")
print()
# 重试死信队列中的事件
platform.subscriptions.retry_dead_letters(
"sub-ticket-completed",
event_ids=[dl.event_id for dl in dead_letters[:10]],
)
print("Retrying 10 dead letter events")
# 清理已过期的死信
platform.subscriptions.purge_dead_letters(
"sub-ticket-completed",
before="2025-01-01T00:00:00Z",
)
#6. 高级模式
#6.1 事件聚合订阅
有时你不想收到每一条变更事件,而是希望在一段时间内汇总后收到一份摘要:
from ontology_sdk.subscription import AggregationWindow
sub_daily_summary = Subscription(
name="ticket-daily-summary",
object_type="Ticket",
event_types=["CREATE", "UPDATE"],
aggregation=AggregationWindow(
window_size="1h", # 每小时汇总一次
aggregate_fields={
"total_created": {"type": "count", "filter": {"event_type": "CREATE"}},
"total_resolved": {
"type": "count",
"filter": {
"event_type": "UPDATE",
"changes.status.after": "COMPLETED",
},
},
"avg_resolution_time_hours": {
"type": "avg",
"field": "resolution_time_hours",
"filter": {"event_type": "UPDATE", "changes.status.after": "COMPLETED"},
},
},
),
target=WebhookTarget(
url="https://dashboard.example.com/api/metrics",
method="POST",
),
)
#6.2 跨对象关联订阅
订阅不局限于单个对象类型,可以通过关联关系设置跨对象的联动订阅:
# 当"项目"下的任何"工单"完成时,检查是否所有工单都完成了
sub_project_completion = Subscription(
name="project-all-tickets-completed",
object_type="Ticket",
event_types=["UPDATE"],
filter=EventFilter(
conditions=[
FilterCondition(field="status", operator="eq", value="COMPLETED", apply_to="after"),
],
),
# 关联检查:触发时检查关联的 Project 下是否所有 Ticket 都完成
correlation=CorrelationCheck(
traverse_link="belongsToProject",
target_object_type="Project",
condition="ALL_LINKED_MATCH",
linked_type="Ticket",
linked_filter=FilterCondition(field="status", operator="eq", value="COMPLETED"),
# 只有当 Project 下所有 Ticket 都是 COMPLETED 时才触发
),
target=InternalActionTarget(
action_type="MarkProjectCompleted",
parameter_mapping={
"project_rid": "${correlation.target_rid}",
},
),
)
#6.3 条件去重
防止短时间内重复触发:
sub_with_dedup = Subscription(
name="alert-dedup",
object_type="MonitoringAlert",
event_types=["CREATE"],
deduplication=DeduplicationConfig(
key_expression="${event.changes.alert_type.after}:${event.changes.host.after}",
window_seconds=300, # 5 分钟内相同 key 只触发一次
strategy="FIRST", # 只保留第一条,或 LAST 保留最后一条
),
target=WebhookTarget(url="https://alert-handler.internal/handle"),
)
#7. 调试与排错
#7.1 事件回放测试
在开发阶段,你可以使用事件回放功能验证订阅逻辑:
# 模拟事件投递(dry-run 模式)
test_event = OntologyEvent(
event_type="UPDATE",
object_type="Ticket",
object_rid="ri.ticket.123",
changes={
"status": Change(before="IN_PROGRESS", after="COMPLETED"),
"resolution_time_hours": Change(before=None, after=4.5),
},
actor="test-user",
)
result = platform.subscriptions.test_delivery(
subscription_id="sub-ticket-completed",
event=test_event,
dry_run=True, # 不实际投递,只返回是否会匹配
)
print(f"Would match: {result.matched}")
print(f"Would deliver to: {result.target_description}")
print(f"Payload preview: {result.payload_preview}")
#7.2 查看投递日志
# 获取最近的投递日志
logs = platform.subscriptions.get_delivery_logs(
"sub-ticket-completed",
limit=20,
status="FAILED", # 只看失败的
)
for log in logs:
print(f"Event: {log.event_id}")
print(f" Time: {log.timestamp}")
print(f" Status: {log.status}")
print(f" Response code: {log.response_code}")
print(f" Response body: {log.response_body[:200]}")
print(f" Duration: {log.duration_ms}ms")
print()
#7.3 常见问题排查
| 问题 | 原因 | 解决方案 |
|---|---|---|
| 订阅不触发 | 过滤条件过严 | 使用 test_delivery 验证过滤逻辑 |
| 重复触发 | 未配置去重 | 添加 deduplication 配置 |
| 延迟高 | 消费端处理慢 | 检查 Webhook 响应时间,增加并发 |
| 事件丢失 | 未确认(ACK) | 确保 gRPC 流式客户端正确 ACK |
| 事件乱序 | 跨分区消费 | 确保同一对象事件在同一分区 |
#8. 生产环境最佳实践
#8.1 订阅设计原则
- 单一职责:每个订阅只做一件事,避免"万能订阅"
- 幂等消费:消费端必须支持幂等,因为事件可能重复投递
- 快速返回:Webhook 处理应在 5 秒内返回,长任务异步化
- 监控 Lag:关注订阅的消费滞后(Lag),告警阈值建议设为 1000
- 优雅降级:消费端不可用时不应影响事件生产
#8.2 性能调优
# 高吞吐场景的订阅配置
high_throughput_sub = Subscription(
name="high-volume-processor",
object_type="SensorReading",
event_types=["CREATE"],
performance=PerformanceConfig(
parallelism=8, # 并行消费线程数
batch_size=100, # 批量投递大小
batch_timeout_ms=1000, # 批量超时(达到 size 或 timeout 就投递)
max_in_flight=1000, # 最大在途消息数
),
target=KafkaTarget(
bootstrap_servers="kafka:9092",
topic="sensor.events.processed",
compression="lz4",
),
)
#8.3 安全建议
- 签名验证:Webhook 投递携带 HMAC 签名,消费端应验证签名
- TLS 加密:生产环境所有 Webhook 和 gRPC 连接必须使用 TLS
- 最小权限:订阅的 Token 只授予读取特定 ObjectType 事件的权限
- 审计日志:所有订阅的创建、修改、删除操作都会记录到审计日志
# Webhook 签名验证示例(消费端)
import hmac
import hashlib
def verify_webhook(request, secret: str) -> bool:
signature = request.headers.get("X-Onto-Signature")
if not signature:
return False
expected = hmac.new(
secret.encode(),
request.body,
hashlib.sha256,
).hexdigest()
return hmac.compare_digest(f"sha256={expected}", signature)
#总结
本教程覆盖了 coomia-dip 事件订阅系统的核心功能:
- 事件模型:理解 CREATE/UPDATE/DELETE/LINK/UNLINK 五种事件类型
- 订阅创建:通过 SDK 声明式创建订阅,配置过滤条件和投递目标
- gRPC 流式订阅:实现低延迟的实时事件消费,包含重连机制
- 实战场景:ITSM 多系统联动的完整实现
- 高级模式:事件聚合、跨对象关联、条件去重
- 运维管理:状态监控、死信队列、暂停恢复
事件订阅是构建响应式系统的基石。在下一篇教程中,我们将探讨如何通过 gRPC 服务扩展平台的能力边界。
下一篇:[S12-14] gRPC 自定义服务开发指南 上一篇:[S12-12] 数据导入最佳实践