Java推送调用流程如何统一

wen java案例 31

本文目录导读:

Java推送调用流程如何统一

  1. 核心设计原则
  2. 标准推送流程架构
  3. 统一消息模型定义
  4. 核心接口与实现
  5. 统一结果与回调处理
  6. 完整调用示例
  7. 最佳实践与注意事项

在Java中统一推送调用流程,核心目标是屏蔽底层差异(如不同厂商SDK:极光、个推、华为、小米、苹果APNs/FCM等),实现业务逻辑与推送技术解耦,以下是一套经过实践验证的架构方案,涵盖从入口到最终送达的完整标准化流程。

核心设计原则

  1. 策略模式:封装不同推送渠道的实现
  2. 责任链模式:处理消息预处理、过滤、路由
  3. 统一消息模型:定义平台无关的推送消息结构
  4. 异步处理:利用消息队列解耦耗时操作

标准推送流程架构

业务系统 → 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)
  • 异步解耦(不影响主业务流程响应时间)

抱歉,评论功能暂时关闭!