本文目录导读:

在Java中统一推送调用流程,核心目标是屏蔽底层差异(如不同厂商SDK:极光、个推、华为、小米、苹果APNs/FCM等),实现业务逻辑与推送技术解耦,以下是一套经过实践验证的架构方案,涵盖从入口到最终送达的完整标准化流程。
核心设计原则
- 策略模式:封装不同推送渠道的实现
- 责任链模式:处理消息预处理、过滤、路由
- 统一消息模型:定义平台无关的推送消息结构
- 异步处理:利用消息队列解耦耗时操作
标准推送流程架构
业务系统 → PushService(统一入口) → PushMessageBuilder → PushFilterChain → PushStrategyRouter → 具体的PushChannel → 回调处理
↓
APNs / FCM / 华为 / 小米 / 个推...
统一消息模型定义
@Data
@Builder
public class UnifiedPushMessage {
// 基础信息
private String pushId; // 推送唯一标识
private String appId; // 应用ID
private PushPlatform platform; // 目标平台枚举
private List<String> targetIds; // 目标用户/设备ID
// 消息内容
private String title;
private String body;
private Map<String, String> extras; // 自定义参数
// 推送策略
private PushPriority priority; // 优先级:HIGH, NORMAL, LOW
private Boolean isOffline; // 是否保留离线消息
private Integer offlineDuration; // 离线保留时长(秒)
private String channelId; // Android渠道ID
private PushUrlType urlType; // 点击跳转类型(APP内页/H5/原生)
private String url; // 跳转地址
// 定时/循环
private LocalDateTime scheduleTime;
private String cronExpression;
}
核心接口与实现
统一推送入口服务
public interface PushService {
PushResult send(UnifiedPushMessage message);
PushResult sendAsync(UnifiedPushMessage message);
PushResult cancel(String pushId);
}
@Service
public class PushServiceImpl implements PushService {
@Autowired
private PushMessageBuilder messageBuilder;
@Autowired
private PushFilterChain filterChain;
@Autowired
private PushStrategyRouter strategyRouter;
@Override
public PushResult send(UnifiedPushMessage message) {
// 1. 构建完整消息
UnifiedPushMessage builtMsg = messageBuilder.build(message);
// 2. 执行过滤链(参数校验、频率限制、内容安全、去重)
boolean passed = filterChain.filter(builtMsg);
if (!passed) {
return PushResult.builder()
.success(false)
.errorCode("FILTER_BLOCKED")
.build();
}
// 3. 路由到具体的推送渠道
PushChannel channel = strategyRouter.route(builtMsg.getPlatform());
// 4. 执行推送(支持重试机制)
return pushWithRetry(channel, builtMsg, 3);
}
}
推送渠道策略路由
@Component
public class PushStrategyRouter {
private final Map<PushPlatform, PushChannel> channelMap = new ConcurrentHashMap<>();
@PostConstruct
public void init() {
channelMap.put(PushPlatform.APNS, new ApnsChannel());
channelMap.put(PushPlatform.FCM, new FcmChannel());
channelMap.put(PushPlatform.HUAWEI, new HuaweiChannel());
channelMap.put(PushPlatform.XIAOMI, new XiaomiChannel());
channelMap.put(PushPlatform.JPUSH, new JPushChannel());
channelMap.put(PushPlatform.GETUI, new GetuiChannel());
}
public PushChannel route(PushPlatform platform) {
PushChannel channel = channelMap.get(platform);
if (channel == null) {
throw new IllegalArgumentException("Unsupported platform: " + platform);
}
return channel;
}
}
统一推送渠道接口
public interface PushChannel {
/**
* 执行单次推送
* @return 推送结果(原始响应)
*/
PushChannelResponse push(UnifiedPushMessage message);
/**
* 查询推送状态
*/
PushChannelResponse queryStatus(String messageId);
/**
* 获取渠道支持的平台
*/
PushPlatform supportedPlatform();
}
// 以华为推送为例
@Component
public class HuaweiChannel implements PushChannel {
private final HuaweiPushClient client; // 华为SDK客户端
@Override
public PushChannelResponse push(UnifiedPushMessage message) {
try {
// 1. 转换消息格式
HwPushMsg hwMsg = convertToHuaweiFormat(message);
// 2. 调用华为SDK
HwPushResponse response = client.sendPush(hwMsg);
// 3. 返回统一结果
return PushChannelResponse.builder()
.success(response.isSuccess())
.messageId(response.getMsgId())
.originalResponse(response)
.build();
} catch (Exception e) {
return PushChannelResponse.builder()
.success(false)
.errorMessage(e.getMessage())
.build();
}
}
private HwPushMsg convertToHuaweiFormat(UnifiedPushMessage msg) {
return HwPushMsg.builder()
.title(msg.getTitle())
.body(msg.getBody())
.targetUserIds(msg.getTargetIds())
.clickAction(HwClickAction.builder()
.type(msg.getUrlType().getCode())
.url(msg.getUrl())
.build())
.build();
}
}
消息过滤链
public interface PushFilter {
boolean filter(UnifiedPushMessage message);
}
@Component
public class PushFilterChain {
private List<PushFilter> filters;
@PostConstruct
public void init() {
filters = Arrays.asList(
new ParameterValidatorFilter(), // 参数校验
new FrequencyLimitFilter(), // 频率限制(同用户1分钟最多5条)
new ContentSecurityFilter(), // 内容安全(敏感词过滤)
new DeduplicationFilter(), // 重复消息检测(相同内容+目标+5分钟内去重)
new BlacklistFilter() // 用户黑名单
);
}
public boolean filter(UnifiedPushMessage message) {
for (PushFilter filter : filters) {
if (!filter.filter(message)) {
log.warn("Push blocked by filter: {}", filter.getClass().getSimpleName());
return false;
}
}
return true;
}
}
异步处理支持(可选但推荐)
@Service
public class AsyncPushProcessor {
@Autowired
private PushServiceImpl pushService;
@Autowired
private RabbitTemplate rabbitTemplate; // 或 KafkaTemplate
// 生产端:异步发送到消息队列
public void sendAsync(UnifiedPushMessage message) {
rabbitTemplate.convertAndSend("push.exchange",
"push.route." + message.getPlatform().name().toLowerCase(),
message);
}
// 消费端:处理消息队列中的推送请求
@RabbitListener(queues = "push.queue.huawei")
public void handleHuaweiPush(UnifiedPushMessage message) {
pushService.send(message);
}
}
统一结果与回调处理
统一推送结果
@Data
@Builder
public class PushResult {
private boolean success;
private String pushId;
private Map<PushPlatform, PlatformResult> platformResults; // 各渠道详细结果
private String errorCode;
private String errorMessage;
private long costTimeMs;
}
@Data
@Builder
public class PlatformResult {
private boolean success;
private String messageId; // 渠道返回的消息ID
private int successCount; // 成功数
private int failCount; // 失败数
private List<String> failTokens; // 失败设备token
private String errorCode;
private String errorMessage;
}
回调处理机制(用于获取送达/点击数据)
@RestController
@RequestMapping("/push/callback")
public class PushCallbackController {
@PostMapping("/{platform}")
public ResponseEntity<String> receiveCallback(
@PathVariable String platform,
@RequestBody CallbackPayload payload) {
// 1. 验证签名
// 2. 解析回调数据
PushCallbackData data = parseCallback(platform, payload);
// 3. 更新数据库中的推送记录状态
pushRecordService.updateStatus(data);
// 4. 触发业务告警(送达率低时)
if (data.getDeliveryRate() < 0.9) {
alertService.sendAlert("推送送达率异常:" + data.getPushId());
}
return ResponseEntity.ok("OK");
}
}
完整调用示例
@RestController
@RequestMapping("/push")
public class PushController {
@Autowired
private PushService pushService;
@PostMapping("/send")
public PushResult sendPush(@RequestBody PushRequest request) {
UnifiedPushMessage message = UnifiedPushMessage.builder()
.pushId(UUID.randomUUID().toString())
.appId("com.example.app")
.platform(request.getPlatform()) // 从请求中获取
.targetIds(request.getTargetIds())
.title(request.getTitle())
.body(request.getBody())
.extras(request.getExtras())
.priority(PushPriority.HIGH)
.isOffline(true)
.offlineDuration(86400) // 1天
.build();
return pushService.send(message);
}
}
最佳实践与注意事项
| 关键点 | 说明 |
|---|---|
| Token管理 | 统一维护设备Token与用户关系,定期清理无效Token |
| 降级策略 | 主渠道失败自动切换到备用渠道(如华为→个推→极光) |
| 限流保护 | 全局限流+用户级限流+渠道级限流 |
| 全链路追踪 | 从业务请求到推送送达,使用TraceId串联所有日志 |
| 测试沙盒 | 为每个渠道准备测试环境,支持黑白名单测试 |
| 配置中心 | 渠道密钥、限流参数、过滤规则通过配置中心动态下发 |
通过这样的架构设计,你的Java推送系统可以实现:
- 新增一个推送渠道仅需实现PushChannel接口(其余代码零改动)
- 统一监控与告警(基于PushResult)
- 灵活的业务过滤(通过修改Filter Chain)
- 异步解耦(不影响主业务流程响应时间)