返回博客

事件订阅与通知:让平台主动告诉你

在前面的教程中,我们掌握了如何主动查询数据、执行 Action、触发规则。但在真实业务场景中,"被动等通知"往往比"主动去轮询"更高效——当订单状态变化时自动通知下游系统,当指标异常时即时推送告警,当审批完成时触发后续流程。

Coomia发布于 2026年1月23日16 分钟阅读
分享本文Twitter / X

系列:S12 开发者教程 · 第 13 篇 | 难度:中级 | 阅读时间:15 分钟

事件订阅与通知:让平台主动告诉你

#引言

在前面的教程中,我们掌握了如何主动查询数据、执行 Action、触发规则。但在真实业务场景中,"被动等通知"往往比"主动去轮询"更高效——当订单状态变化时自动通知下游系统,当指标异常时即时推送告警,当审批完成时触发后续流程。

coomia-dip 的事件订阅系统正是为此设计的。它基于本体模型中的对象变更事件,提供了声明式的订阅机制。你只需要告诉平台"我关心什么",平台就会在事件发生时主动将消息推送给你。本教程将带你从零搭建一套完整的事件订阅与通知体系。

#1. 理解 coomia-dip 事件模型

#1.1 事件的本质

在 coomia-dip 中,一切变更都是事件。当你通过 Action 修改一个对象的属性、创建一个新的关联关系、或者删除一个对象时,平台内部都会产生对应的变更事件。这些事件构成了一条不间断的事件流(Event Stream),订阅系统的工作就是从这条流中筛选出你关心的事件并投递给你。

Python
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):每个事件至少会被投递一次到每个活跃订阅
  • 事件有序性:同一个对象的事件保证按照发生顺序投递
  • 持久化存储:事件在投递前会持久化到事件存储,即使系统重启也不会丢失
  • 重放能力:支持从历史某个时间点开始重放事件
Python
# 事件存储的持久化机制
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 创建订阅

让我们从一个实际场景开始:你需要在"工单"对象的状态变为"已完成"时收到通知。

Python
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 支持丰富的过滤表达式:

Python
# 示例 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 支持多种投递目标:

Python
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 流式订阅接口:

PROTOBUF
// 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 客户端实现

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 带重连的健壮客户端

生产环境中,网络断开和服务重启是常态。以下是一个带自动重连的健壮实现:

Python
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),需要实现以下联动:

  1. 工单创建 → 通知 Slack 频道
  2. 工单分配 → 发送邮件给负责人
  3. 工单完成 → 触发满意度调查
  4. SLA 即将到期 → 升级告警

#4.2 实现方案

Python
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 查看订阅状态

Python
# 列出所有订阅
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 暂停与恢复

Python
# 暂停订阅(维护时使用)
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)

投递失败超过重试次数的事件不会丢失,而是进入死信队列:

Python
# 查看死信队列中的事件
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 事件聚合订阅

有时你不想收到每一条变更事件,而是希望在一段时间内汇总后收到一份摘要:

Python
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 跨对象关联订阅

订阅不局限于单个对象类型,可以通过关联关系设置跨对象的联动订阅:

Python
# 当"项目"下的任何"工单"完成时,检查是否所有工单都完成了
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 条件去重

防止短时间内重复触发:

Python
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 事件回放测试

在开发阶段,你可以使用事件回放功能验证订阅逻辑:

Python
# 模拟事件投递(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 查看投递日志

Python
# 获取最近的投递日志
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 订阅设计原则

  1. 单一职责:每个订阅只做一件事,避免"万能订阅"
  2. 幂等消费:消费端必须支持幂等,因为事件可能重复投递
  3. 快速返回:Webhook 处理应在 5 秒内返回,长任务异步化
  4. 监控 Lag:关注订阅的消费滞后(Lag),告警阈值建议设为 1000
  5. 优雅降级:消费端不可用时不应影响事件生产

#8.2 性能调优

Python
# 高吞吐场景的订阅配置
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 事件的权限
  • 审计日志:所有订阅的创建、修改、删除操作都会记录到审计日志
Python
# 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] 数据导入最佳实践