首页 / 工程规范 / 消息与定时任务
消息与定时任务
NATS/RocketMQ/Redis Stream 消息总线、Outbox 与 XXL-JOB 处理器清单。
消息总线实现
| 实现 | 使用位置 | 排障入口 |
|---|---|---|
| 本地总线 | IoT/GPS/Alarm 的 messagebus local 实现 | 单体开发默认先看本地配置与注册日志 |
| Redis Stream | IoT core 的 Redis message bus | 检查 stream、consumer group 与 ACK 日志 |
| RocketMQ | Notify、IoT 绑定 Saga 等消费者/生产者 | 检查 nameserver、topic、consumer group 与重试 |
| NATS JetStream | IoT、GPS、Alarm 的 NATS messagebus | 检查 stream、subject、durable/group 与 ACK |
| Kafka / RabbitMQ | starter 与 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-v1 | IoT 绑定 + 订阅 Seat Saga | subscription.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:
PayOutboxPublishJob、IotBindingOutboxJob、SubscriptionBindingOutboxJob、GpsReportOutboxPublishJob。 - 幂等约束:绑定任务
task_no/pay_order_no、支付 attempt 幂等键、SIM 回调幂等键、订阅 seat 唯一约束。 - 消费组:多副本必须指定 Queue Group/consumerGroup,避免重复消费。
XXL-JOB 处理器清单
| 模块 | 处理器 |
|---|---|
| iot | iotBindingOutboxJob、iotBindingRetryJob、iotOtaUpgradeJob |
| pay | payNotifyJob、payOrderExpireJob、payOrderSyncJob、payUnknownAttemptRecoveryJob、payRefundSyncJob、payTransferSyncJob、payReconciliationCaseJob、paySettlementReconciliationJob、payProviderEventJob、payProviderResourceSyncJob、payOutboxPublishJob |
| subscription | subscriptionExpireJob、subscriptionBindingOutboxJob |
| sim | simProviderSyncJob、simCallbackProcessJob |
| gps | gpsReportScheduleJob、gpsReportOutboxPublishJob、gpsReportNotificationJob |
| infra | accessLogCleanJob |
| system | demoJob |
执行器 AppName 由 spring.application.name + profile 组成;调度中心地址、Token 与执行器日志路径从 profile 注入。新增消费者或 Job 必须同步配置、幂等/重试策略、监控日志与部署环境。