首页 / 工程规范 / 消息与定时任务

消息总线实现

实现使用位置排障入口
本地总线IoT/GPS/Alarm 的 messagebus local 实现单体开发默认先看本地配置与注册日志
Redis StreamIoT core 的 Redis message bus检查 stream、consumer group 与 ACK 日志
RocketMQNotify、IoT 绑定 Saga 等消费者/生产者检查 nameserver、topic、consumer group 与重试
NATS JetStreamIoT、GPS、Alarm 的 NATS messagebus检查 stream、subject、durable/group 与 ACK
Kafka / RabbitMQstarter 与 profile 配置存在只有找到具体 listener/producer 和启用配置后才判定环境依赖
处理失败时先确认消息是否被 ACK、是否进入重试/Outbox,再查看业务表;不要直接修改业务行伪造成功。

NATS 主题规范

常量统一在 go-share/events/topics/nats.go

  • 设备上下行:device.up/down.<protocol>.<handle>.<productID>.<deviceName>
  • 物模型:device.up/down.thing.{productID}.{deviceName}
  • 网关/子设备:device.up/down.gateway.*
  • OTA / 影子 / 配置 / 日志:device.up/down.{ota|shadow|config|log}.*
  • 状态:device.up.status.connected / disconnected
  • 应用事件:application.device.{productID}.{deviceName}.report.thing.property / .event / .action,以及 .status.connected/disconnected

RocketMQ 关键 Topic

Topic用途Tag(事件类型)
entrax-subscription-binding-saga-v1IoT 绑定 + 订阅 Seat Sagasubscription.seat.reserve.requested / .reserved / .reserve-failed / .payment-required / .paid.activation.requested / .completed / .release.requested / .released / iot.binding.completed / .failed
notify / mail / sms消息任务、邮件、短信发送按 consumer 配置的 selectorExpression

Outbox 与幂等

  • Outbox 模式:业务写库 + Outbox 写入同事务,Job 异步发布消息,保证"业务成功必发消息"。
  • 已知 Outbox Job:PayOutboxPublishJobIotBindingOutboxJobSubscriptionBindingOutboxJobGpsReportOutboxPublishJob
  • 幂等约束:绑定任务 task_no/pay_order_no、支付 attempt 幂等键、SIM 回调幂等键、订阅 seat 唯一约束。
  • 消费组:多副本必须指定 Queue Group/consumerGroup,避免重复消费。

XXL-JOB 处理器清单

模块处理器
iotiotBindingOutboxJobiotBindingRetryJobiotOtaUpgradeJob
paypayNotifyJobpayOrderExpireJobpayOrderSyncJobpayUnknownAttemptRecoveryJobpayRefundSyncJobpayTransferSyncJobpayReconciliationCaseJobpaySettlementReconciliationJobpayProviderEventJobpayProviderResourceSyncJobpayOutboxPublishJob
subscriptionsubscriptionExpireJobsubscriptionBindingOutboxJob
simsimProviderSyncJobsimCallbackProcessJob
gpsgpsReportScheduleJobgpsReportOutboxPublishJobgpsReportNotificationJob
infraaccessLogCleanJob
systemdemoJob

执行器 AppName 由 spring.application.name + profile 组成;调度中心地址、Token 与执行器日志路径从 profile 注入。新增消费者或 Job 必须同步配置、幂等/重试策略、监控日志与部署环境。