跳到主要内容

Event Integration

NE503 的 Event Bus 是设备内部的发布/订阅消息中枢。本页只讲事件协议和接入方式;HTTP 认证、请求体和事件管理接口请看 RESTful API 参考,应用权限请看 应用参考

1. 先选择接入方式

接入方式默认端点适用场景位置
Python SDKunix:///run/aipc/event-bus.sock容器内应用发布或订阅使用 hailo_ipc_sdk.events.EventClient
gRPC同上;由部署环境提供可达端点C++、Go 或自定义服务使用 event.proto
REST/api/v1/events/topics/api/v1/events/publish外部系统查询主题或发布事件需要 API 认证
WebSocket/api/v1/events/stream浏览器或外部系统实时接收需要 API 认证
MQTT 桥接由桥接程序连接上述接口对接已有 MQTT 平台需要自己维护桥接进程

设备内部的 gRPC TCP 地址如果被部署为 loopback,只能由设备本机访问;外部主机不要假设可以直接连接,应使用 REST/WebSocket 或经批准的隧道方案。

2. 主题匹配

主题使用 / 分隔的层级字符串。当前匹配器支持精确匹配、* 单级匹配、** 多级匹配,以及 **/suffix 形式的后缀匹配。

订阅表达式匹配示例不匹配示例
inference/model_a/main完全相同的主题inference/model_b/main
inference/*/maininference/model_a/maininference/model_a/sub/extra
inference/**inference/model_a/main、更深层级主题不属于 inference/ 的主题
**/detections任意前缀下以 detections 结尾的主题.../detections/raw

通配符只解决主题匹配,不会替订阅者解析 payload。高吞吐场景应缩小订阅范围,并根据 queue_sizedrop_old 和消费者处理速度做取舍。

3. Event 消息结构

协议定义在工程源码 platform/event-bus/proto/event.proto

{
"topic": "app/alert",
"timestamp_ns": 1717545600000000000,
"source": "people_counting",
"event_id": "evt-1",
"payload": "{\"type\":\"person_detected\"}",
"payload_type": "json",
"metadata": {
"stream_id": "main"
}
}
字段类型说明
topicstring主题名称
timestamp_nsuint64纳秒时间戳
sourcestring服务名或应用 ID
event_idstring事件 ID;发布响应也会返回 ID
payloadbytesJSON 或其他编码的内容
payload_typestring当前常用 json,也可标记 protobuf
metadatamap of string to string可选元数据

Python SDK 的 EventClient 会把 JSON payload 转为 event.payload 字典;原始 gRPC 客户端需要根据 payload_type 自己解析。发布请求还支持 persistentttl_ms

订阅请求支持 topicsubscriber_idfiltersqueue_sizedrop_old;服务端提供 PublishPublishBatchSubscribeUnsubscribe、主题查询和统计 RPC。

4. 事件来源与主题形态

主题前缀用于快速定位来源,但不是固定的全量枚举:

来源源码中的主题形态说明
App Managerapp/{eventType}应用安装、启动、停止等生命周期事件
Device Controldevice/{eventType}设备控制服务发布的设备事件
AI Runtime 常规结果inference/{model_id}/{stream_id}grpc_service.cpp 的结果发布路径
AI Runtime 自动推理路径inference/{stream_id}auto_infer.cpp 中的另一条发布路径
应用自定义事件由应用自行定义应与 app/{app_id}/... 等命名约定保持一致

AI Runtime 的两条源码路径并不具有相同的段数。订阅第三方推理结果时优先使用设备当前实际出现的主题,并结合 event.metadata 或 payload 判断模型和码流;不要对所有固件版本都硬编码三段式主题。

5. Python SDK:发布与订阅

from hailo_ipc_sdk.events import EventClient


def main():
with EventClient() as events:
events.publish(
"app/people_counting/stats",
{"current_count": 2, "threshold": 10},
metadata={"stream_id": "main"},
)

for event in events.subscribe(
"app/**",
queue_size=100,
drop_old=True,
):
print(event.topic, event.payload, event.source)


if __name__ == "__main__":
main()

EventClient 还提供 publish_batch()on_event()unsubscribe()list_topics()get_topic_info()get_stats()get_topic_stats()。订阅是阻塞迭代器;长时间运行的应用应处理退出信号并调用 close()

5.1 权限清单

应用使用 Event Bus 前,在 app.yaml 中声明主题范围:

permissions:
events:
publish:
- app/people_counting/*
subscribe:
- inference/**

发布和订阅权限分别校验。不要为了省事使用过大的 **,除非应用确实需要接收所有主题。

6. WebSocket:外部实时订阅

事件 WebSocket 路径为:

wss://<设备IP>/api/v1/events/stream

先通过 REST 登录,再按 RESTful API 参考 的认证约定建立连接。服务端在源码中以通配符订阅 Event Bus,再把事件转发给 WebSocket 客户端;因此 URL 上附加的主题过滤不能替代客户端过滤。客户端应根据 topicsource 或 payload 做二次筛选,并处理断线重连和重复事件。

如果业务只需要设备内的应用间通信,优先使用 Unix Socket 上的 SDK/gRPC;WebSocket 更适合浏览器、网关或外部监控系统。

7. MQTT 桥接

MQTT 不是 Event Bus 的原生协议。桥接程序作为 Event Bus 订阅者,再把事件转换成 MQTT 消息:

import json
import os

import paho.mqtt.client as mqtt
from hailo_ipc_sdk.events import EventClient


mqtt_client = mqtt.Client(client_id=os.environ["MQTT_CLIENT_ID"])
mqtt_client.username_pw_set(
os.environ["MQTT_USERNAME"], os.environ["MQTT_PASSWORD"]
)
mqtt_client.connect(os.environ["MQTT_HOST"], int(os.environ.get("MQTT_PORT", "1883")))
mqtt_client.loop_start()

events = EventClient()
try:
for event in events.subscribe("**"):
payload = json.dumps({
"timestamp_ns": event.timestamp_ns,
"source": event.source,
"event_id": event.event_id,
"payload": event.payload,
"metadata": event.metadata,
})
mqtt_client.publish(f"ne503/{event.topic}", payload, qos=1)
finally:
events.close()
mqtt_client.loop_stop()
mqtt_client.disconnect()

生产桥接需要补充 TLS、认证密钥管理、重连、QoS 和重复事件处理。不要把账号密码写入文档示例或镜像源码。

8. 集成落地时的三个决定

  1. 主题范围:按业务订阅最小主题集合,避免无关事件挤占队列。
  2. 消息语义:payload 中保留业务对象、时间戳和必要的幂等键;不要只依赖接收时间。
  3. 失败处理:为 SDK、WebSocket 和 MQTT 消费者定义重连、丢弃旧消息、重复消费和下游不可用时的策略。

9. 相关文档