本文目录导读:

规范Java消息推送流程,需要从健壮性、可追溯性、可靠性三个核心维度出发,一个完善的推送流程通常包含:消息构建 → 消息存储 → 消息分发/路由 → 推送执行 → 回执与重试 → 闭环监控。
以下是一套标准化的规范流程及关键实现要点:
第一阶段:消息标准化与构建
规范点:
- 统一消息体结构: 设计通用的
Message对象,包含以下标准字段:messageId:全局唯一ID(如UUID或雪花算法生成)。bizCode:业务类型码(如ORDER_PAID、SYSTEM_NOTIFY)。userId:目标用户ID。channel:推送渠道列表(如["SMS", "APP_PUSH", "WEBSOCKET"])。title&content。priority:优先级(LOW/HIGH/URGENT)。expireTime:消息过期时间(避免积压推送)。extra:扩展字段(如跳转链接、JSON参数)。
- 序列化与反序列化统一: 推荐使用Protobuf或JSON Schema进行序列化,保证跨语言、跨系统兼容。
第二阶段:消息持久化与防丢失
规范点:
- 强制落库(异步写入): 消息构建后,不能直接发送。
- 步骤1: 先将消息写入数据库(如MySQL的
message_queue表)或消息队列(如Kafka,并开启acks=all)。 - 步骤2: 消息状态初始化为
PENDING(待发送)。
- 步骤1: 先将消息写入数据库(如MySQL的
- 刷盘策略: 对于重要消息(如交易通知),必须保证
消息写入成功后才返回给调用端成功(保证At Least Once,即至少一次送达)。
第三阶段:消息路由与分发
规范点:
- 渠道适配器模式: 定义
ChannelAdapter接口(如SmsAdapter、PushAdapter、WebsocketAdapter),根据消息体中的channel字段动态分发。 - 消息过滤与合并:
- 去重: 根据
messageId或bizCode对短时间内重复请求进行去重。 - 合并: 同一用户1分钟内收到多条同类型通知,合并为一条(如“您有3条未读消息”)。
- 去重: 根据
- 异步投递: 使用线程池或响应式编程(如WebFlux)推送,避免阻塞业务主流程。
第四阶段:推送执行与自适应
规范点:
- 限流与熔断(关键):
- 第三方渠道限流: 针对短信、邮件等有QPS限制的渠道,使用令牌桶或滑动窗口限流(如Guava RateLimiter或Sentinel)。
- 用户级别限流: 防止对单一用户频繁推送(如1分钟内不超过5条)。
- 熔断: 若第三方推送接口连续失败(如响应超过5秒或返回503),触发熔断(断开该渠道连接,降级为备用渠道)。
- 幂等性保证: 推送逻辑必须幂等,通过数据库唯一索引(
messageId+userId+channel)防止重复投递,如果第一次已发送成功,第二次直接返回成功或忽略。
第五阶段:回执处理与补偿机制
规范点:
- 状态更新:
PENDING→SENT(发送成功)或FAILED(发送失败)。- 对于App推送(如APNs/FCM),通过长连接接收到达回执(Delivery Receipt)。
- 补偿重试:
- 立即重试: 网络抖动导致的失败,重试2-3次(指数退避算法,如:1s -> 2s -> 4s)。
- 定时重试: 对于非瞬时错误(如设备Token失效),标记为
INVALID_DEVICE,不再重试;对于系统异常,放入重试队列(延时队列)。 - 死信处理: 重试超过最大次数(如3次)的消息,进入死信队列(Dead Letter Queue,DLQ),人工或脚本介入分析。
- 定时任务兜底: 建立定时任务(如每分钟扫描一次),扫描
status = PENDING且expireTime > now()的消息进行补推(适用于可靠消息场景)。
第六阶段:监控与闭环
规范点:
- 全链路追踪: 使用Trace ID(如整合SkyWalking、Zipkin)记录消息从构建到送达的全部链路耗时。
- 核心Metrics:
- 发送量: 按渠道、业务、用户维度的成功/失败数量。
- 延迟: 消息从创建到送达的P99延迟。
- 到达率: (送达数 / 发送数)× 100%。
- 转化率: 推送后用户点击率(需配合行为埋点)。
- 告警规则:
- 某渠道失败率 > 5%。
- 消息延迟超过阈值(如SMS > 30秒)。
- 死信队列积累超过100条。
核心架构流程图(简化版)
[业务系统]
|
v
[统一API接入层] --> 生成messageId
|
v
[消息持久化] ----> 状态: PENDING
| |
| v
| [异步分发器] ----> 按渠道路由
| |
v v
[消息队列] --> [推送执行器] --> RateLimiter --> [第三方渠道(推送/SMS/WS)]
| |
v v
[状态更新] [回调处理] --> 更新status为SENT/FAILED
| |
v v
[死信队列] <--- [重试队列] <--- [定时任务扫描]
|
v
[人工介入/分析]
常见陷阱与避坑指南
- 避免同步阻塞: 绝对不能直接在HTTP请求中同步调用第三方推送(如短信API),这会导致接口超时。必须异步化。
- 用户退订处理: 推送前必须检查用户是否在
unsubscribe黑名单中(尤其是营销类推送),否则可能被平台封禁(如微信、钉钉)。 - 敏感信息脱敏: 日志中打印消息内容时,必须脱敏(如手机号中间四位 、身份证号)。
message存入MQ或数据库时建议加密存储。 - 通道降级: 设计
ChannelRouter时,若App Push失败(如用户离线),应自动降级为短信或静默通知(需在消息体预先定义降级策略)。
一个规范的Java消息推送流程,本质上是一个异步、可靠、可观测的消息状态机。
- 核心:消息必须先落库再发送。
- 关键:必须有限流、重试、死信机制。
- 保障:通过全链路Trace和Metrics实现可观测。
遵循上述流程,可以有效避免消息丢失、重复推送、系统雪崩等常见问题。