首页 / 数据流转 / 业务事件流
业务事件流
NATS / RocketMQ 事件总表:生产者、消费者、Topic/Tag 与消费组规范。
事件总线总览
| 总线 | 典型用途 | 可靠性 |
|---|---|---|
| NATS / JetStream | 设备上下行、物模型事件、告警事件(跨语言主干) | JetStream 持久化 + Queue Group |
| RocketMQ | 绑定 Saga、通知/短信/邮件任务 | 重试 + 事务消息 |
| Redis Stream | IoT 本地消息总线(可选实现) | Stream + Consumer Group |
| 本地总线 | 单体开发默认(IoT/GPS/Alarm) | 进程内 |
NATS 主题与消费方
| 主题 | 生产者 | 消费者 | 落点 |
|---|---|---|---|
device.up.thing.{productID}.{deviceName} | dgsvr/gpscodecsvr | dmsvr、geosvr、Java | 物模型处理 |
GPS_RAW_V1(Stream) | dmsvr/编解码层 | geosvr | TDengine 原始轨迹 |
GPS_QUALITY_V1(Stream) | geosvr | tripsvr 等 | 质量清洗后轨迹 |
DMS_DEVICE_THING_UP_V1(Stream) | dgsvr | dmsvr | 设备物模型上行 |
device.up.status.connected/disconnected | dgsvr | Java iot 等 | 在线状态业务 |
application.device.*.report.thing.property | Go 数据面 | Java(iots、gps、alarm) | 业务规则与告警 |
device.down.thing/ota/config/shadow.* | Java/Go 业务 | dmsvr/dgsvr → 设备 | 下行指令 |
RocketMQ Topic 与消费方
| Topic | Tag(事件) | 消费者组 | 消费者 |
|---|---|---|---|
entrax-subscription-binding-saga-v1 | seat.reserve.requested 等 | subscription-binding-saga-subscription-v1 | SubscriptionBindingSagaListener(subscription) |
seat.reserved / paid.activation-completed 等 | subscription-binding-saga-iot-v1 | IotBindingSagaListener(iot) | |
| notify 消息事件 | 站内信/推送创建 | 按模块配置 | NotifyMessageTaskConsumer 等 |
| 短信/邮件 | 发送任务 | 按模块配置 | SmsSendConsumer、MailSendConsumer |
事件命名与消费规范
- NATS 主题常量统一在
go-share/events/topics/nats.go,禁止硬编码字符串。 - RocketMQ Saga 事件类型统一在
SubscriptionBindingSagaConstants(SCHEMA_VERSION=1)。 - 多副本消费必须指定 Queue Group / consumerGroup,避免重复消费。
- 不可丢失消息(告警、订单状态)走 JetStream 持久化 + 消费端幂等。
- 消费失败先确认 ACK/重试/Outbox,再查业务表;不伪造业务行。