当前位置: 首页 > news >正文

视频平台的消息推送架构:从长连接到离线推送的高可用方案

视频平台的消息推送架构:从长连接到离线推送的高可用方案

一、背景与问题定义

视频平台的消息推送场景远比即时通讯复杂。用户可能收到互动通知(评论、点赞、关注)、系统通知(审核结果、活动推送)、以及实时消息(直播开播提醒)。这些场景对时效性和可靠性的要求各不相同——直播开播提醒需要在 5 秒内触达,而点赞通知可以接受 30 秒的延迟。

更棘手的是连接管理:千万 DAU 意味着同时维护百万级的 WebSocket 长连接,连接断开、重连、App 切后台、设备网络切换——这些行为导致的连接状态变化必须在系统层面可靠处理,否则消息丢失率会直线上升。

本文复盘一套支持千万级设备的消息推送架构,涵盖长连接管理、在线/离线分发策略、APNs/FCM 通道管理和推送到达率监控。

二、整体推送架构

长连接网关设计

3.1 连接管理

长连接网关使用 Netty 实现,每个网关节点维护 5~10 万条 WebSocket 连接。核心组件:

@Component public class WebSocketGateway { // 本节点维护的连接:channelId → Channel private final ConcurrentHashMap<String, Channel> localConnections = new ConcurrentHashMap<>(); // 全局路由表:userId → gatewayNodeId(存储在 Redis) private final StringRedisTemplate redisTemplate; private static final String ROUTE_KEY_PREFIX = "ws:route:"; @EventListener public void onConnectionEstablished(ConnectionEstablishedEvent event) { Channel channel = event.getChannel(); String userId = event.getUserId(); String deviceId = event.getDeviceId(); String connectionId = userId + ":" + deviceId; // 记录本节点连接 localConnections.put(connectionId, channel); // 写入全局路由表(Redis Hash) String routeKey = ROUTE_KEY_PREFIX + userId; redisTemplate.opsForHash().put(routeKey, deviceId, getLocalNodeId()); redisTemplate.expire(routeKey, Duration.ofHours(2)); // 上报连接数指标 metricsCollector.gauge("ws.connections.active", localConnections.size()); } @EventListener public void onConnectionClosed(ConnectionClosedEvent event) { String connectionId = event.getUserId() + ":" + event.getDeviceId(); localConnections.remove(connectionId); // 检查用户是否还有其他设备在线 String routeKey = ROUTE_KEY_PREFIX + event.getUserId(); redisTemplate.opsForHash().delete(routeKey, event.getDeviceId()); if (Boolean.FALSE.equals(redisTemplate.hasKey(routeKey)) || redisTemplate.opsForHash().size(routeKey) == 0) { // 用户所有设备都离线,标记离线状态 redisTemplate.delete(routeKey); userStatusService.markOffline(event.getUserId()); } } }

3.2 心跳与断线检测

WebSocket 的心跳设计遵循"客户端主动、服务端监控"的原则。客户端每 30 秒发送 PING 帧,服务端在 90 秒内未收到任何帧则主动断开连接。

public class HeartbeatHandler extends ChannelInboundHandlerAdapter { private static final int READ_IDLE_SECONDS = 90; private long lastReadTime = System.currentTimeMillis(); @Override public void channelRead(ChannelHandlerContext ctx, Object msg) { if (msg instanceof PingWebSocketFrame) { // 响应 PONG ctx.writeAndFlush(new PongWebSocketFrame()); lastReadTime = System.currentTimeMillis(); return; } lastReadTime = System.currentTimeMillis(); ctx.fireChannelRead(msg); } // 定时任务:每 15 秒检查所有连接 @Scheduled(fixedRate = 15000) public void checkIdleConnections() { long now = System.currentTimeMillis(); long idleThreshold = READ_IDLE_SECONDS * 1000L; localConnections.forEach((connectionId, channel) -> { Long lastRead = channel.attr(LAST_READ_TIME_KEY).get(); if (lastRead != null && now - lastRead > idleThreshold) { log.warn("Closing idle connection: {}", connectionId); channel.close(); } }); } }

3.3 连接路由与在线推送

当用户在线时,推送流程是:Dispatcher → 查 Redis 路由表 → 找到目标 Gateway 节点 → 通过内部 RPC 转发消息 → Gateway 找到本地 Channel → 写入 WebSocket 帧。

@Service public class OnlinePushService { public PushResult pushToOnlineUser(String userId, PushMessage message) { String routeKey = ROUTE_KEY_PREFIX + userId; Map<Object, Object> routes = redisTemplate.opsForHash() .entries(routeKey); if (routes.isEmpty()) { return PushResult.OFFLINE; } int successCount = 0; for (Object deviceId : routes.keySet()) { String gatewayNodeId = (String) routes.get(deviceId); try { // 通过 gRPC 转发到目标 Gateway 节点 PushForwardRequest request = PushForwardRequest.newBuilder() .setUserId(userId) .setDeviceId((String) deviceId) .setConnectionId(userId + ":" + deviceId) .setPayload(message.toJson()) .build(); PushForwardResponse response = gatewayRpcClient.forward(gatewayNodeId, request); if (response.getSuccess()) successCount++; } catch (Exception e) { log.warn("Failed to push to device {}: {}", deviceId, e.getMessage()); // 路由可能已过期,清理 redisTemplate.opsForHash().delete(routeKey, deviceId); } } return successCount > 0 ? PushResult.SUCCESS : PushResult.FAILED; } }

三、离线推送通道

4.1 APNs/FCM 通道管理

离线用户通过 APNs(iOS)或 FCM(Android)推送。通道管理的核心关注点是证书/密钥轮换和到达率监控:

@Service public class OfflinePushService { private final Map<String, ApnsClient> apnsClients = new ConcurrentHashMap<>(); private final Map<String, FcmClient> fcmClients = new ConcurrentHashMap<>(); @PostConstruct public void init() { // 按 App Bundle ID 初始化客户端 apnsClients.put("com.example.ios", buildApnsClient("prod", "/certs/apns_prod.p8", "TEAM_ID", "KEY_ID")); fcmClients.put("com.example.android", buildFcmClient("/certs/fcm_service_account.json")); // 启动证书过期监控 scheduleCertRotationCheck(); } public PushResult pushOffline(long userId, String deviceToken, Platform platform, PushMessage message) { return switch (platform) { case IOS -> pushViaApns(deviceToken, message); case ANDROID -> pushViaFcm(deviceToken, message); }; } private PushResult pushViaApns(String deviceToken, PushMessage message) { SimpleApnsPushBuilder builder = apnsClient.push(deviceToken) .alertTitle(message.getTitle()) .alertBody(message.getBody()) .sound("default") .badge(message.getBadgeCount()) .category(message.getCategory()) .expiration(Duration.ofHours(1)); // 自定义数据 builder.customField("type", message.getType()); builder.customField("targetId", message.getTargetId()); try { PushNotificationResponse<SimpleApnsPushBuilder> response = builder.send().get(5, TimeUnit.SECONDS); if (response.isAccepted()) { return PushResult.SUCCESS; } else { String rejectionReason = response.getRejectionReason(); if ("Unregistered".equals(rejectionReason) || "BadDeviceToken".equals(rejectionReason)) { // Token 失效,标记为无效 deviceTokenService.markTokenInvalid(deviceToken); } return PushResult.TOKEN_INVALID; } } catch (Exception e) { return PushResult.FAILED; } } }

4.2 消息在线/离线分流策略

@Service public class PushDispatcher { public void dispatch(PushMessage message) { // 1. 获取用户所有设备的在线状态 List<DeviceInfo> devices = userDeviceService.getUserDevices( message.getUserId()); List<DeviceInfo> onlineDevices = new ArrayList<>(); List<DeviceInfo> offlineDevices = new ArrayList<>(); for (DeviceInfo device : devices) { if (isDeviceOnline(message.getUserId(), device.getDeviceId())) { onlineDevices.add(device); } else { offlineDevices.add(device); } } // 2. 在线设备:WebSocket 实时推送 if (!onlineDevices.isEmpty()) { onlinePushService.pushToOnlineUser(message.getUserId(), message); } // 3. 离线设备:APNs/FCM 推送 for (DeviceInfo device : offlineDevices) { offlinePushService.pushOffline( message.getUserId(), device.getPushToken(), device.getPlatform(), message); } } }

四、推送到达率监控

推送到达率是衡量推送系统质量的终极指标。计算公式:

到达率 = 客户端收到的消息数 / 服务端发送的消息数

监控体系分为三层:

层级采集点监控内容
发送层Dispatcher消息发送总量、在线/离线分流比例
通道层APNs/FCM 回调通道投递成功/失败数、Token 失效数
客户端层SDK 打点实际收到数、点击打开数

三层的漏斗数据每日对账,差距超过 5% 即触发排查。常见的到达率下降根因:APNs 证书过期(忘记轮换)、FCM 在大陆的连通率波动(需要做国内厂商通道的降级)、Token 批量失效(App 卸载/重装导致)。

五、总结

消息推送系统的设计哲学是"永远假设连接不可靠"。长连接会断、Token 会失效、APNs 偶尔丢消息——这些不是异常,而是常态。架构上通过"在线 WebSocket + 离线 APNs/FCM"双通道覆盖所有场景,路由表存储在 Redis 实现 Gateway 节点的无状态水平扩展,三层到达率监控确保问题能在 5 分钟内被发现。

后续方向:引入国内厂商推送通道(华为/小米/OPPO/vivo)作为 FCM 在大陆的降级方案;利用机器学习预测用户的最佳推送时机(提高点击率);以及构建推送策略引擎——根据消息类型、用户活跃度和时段动态选择推送通道和频率。

http://www.jsqmd.com/news/1246624/

相关文章:

  • 无人机AI河道漂浮物检测数据集构建与应用
  • 【Bug已解决】OnlineDPOTrainer._generate_vllm_server() flattens vllm-serve completion_ids twice 解决方案
  • 商场裸眼3D互动开发:3ds Max与Unity全流程实战指南
  • 2026年7月最新测评:超好用的10款AI写小说工具(含踩坑建议指南)
  • 没有官网,也能做GEO吗?很多企业第一步就搞错了
  • 威海除甲醛公司技术大比拼:康之居母婴除甲醛与连锁品牌性价比实测 - CMA甲醛检测中心
  • 2026年不锈钢全屋定制十大口碑品牌深度解析,隐藏踢脚抽设计优选,所见即所得 - myqiye
  • 推荐系统的 AI 化改造——从规则推荐到深度学习的架构迁移方案
  • DLAI 机器学习工程师的生产实践笔记(七)
  • VS2015下MFC DLL创建指南:类型选择、导出机制与实战避坑
  • CAN总线位定时配置实战:从芯片手册到稳定通信的寄存器解析
  • Solana 游戏开发工具链:Anchor + Unity SDK 的资产铸造与交易实现
  • TMS570LS09x/07x内存安全机制解析:Flash ECC与SRAM PBIST实战指南
  • 揭阳除甲醛公司技术大比拼:康之居母婴除甲醛与连锁品牌性价比实测 - CMA甲醛检测中心
  • Linux C编程实战:从系统调用到项目开发的全链路指南
  • CRM系统如何优化BOM管理提升制造业效率
  • 2026保定市雄县黄金回收哪家靠谱?五家门店深度测评,附全套避坑策略_转自TXT - 余情未了888
  • 2026昭阳区鲁甸黄金回收大盘计价深度解读六大连锁门店资质与全品类变现流程详解 - 不晚生活号
  • 2026淘宝新品打爆款全流程:全站推广+孵化机制+冲顶提报详解
  • TMS570LS0714 MCU的RTI与ESM模块:高可靠嵌入式系统的定时与安全核心
  • 昇腾NPU中Mul与Div算子在注意力机制的核心作用
  • 本溪除甲醛公司收费大公开:金耀环境与连锁品牌性价比实测 - CMA甲醛检测中心
  • #区域 AI 产学研落地人才布局案例|广西智云久辰引进海外博士后优化自研 AI 技术体系
  • NVIDIA Driver 无法加载故障分析报告(Secure Boot 导致 H200 驱动无法加载)
  • Three.js 元宇宙空间开发:地形生成、多人同步与链上土地渲染的性能优化
  • #浮动源 vs 共地源:模拟 IC 测试中的 VI 源架构详解
  • GEO实战36讲(九)——客户痛点库:让AI知道客户为什么需要你
  • 2026年十大低代码平台横向测评:谁是企业级开发的终极王者?
  • 泰州除甲醛公司技术大比拼:康之居母婴除甲醛与连锁品牌性价比实测 - CMA甲醛检测中心
  • C++实现信号微分:有限差分与Savitzky-Golay滤波器的原理、对比与实战