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

发布时间:2026/7/23 8:19:03

视频平台的消息推送架构:从长连接到离线推送的高可用方案 视频平台的消息推送架构从长连接到离线推送的高可用方案一、背景与问题定义视频平台的消息推送场景远比即时通讯复杂。用户可能收到互动通知评论、点赞、关注、系统通知审核结果、活动推送、以及实时消息直播开播提醒。这些场景对时效性和可靠性的要求各不相同——直播开播提醒需要在 5 秒内触达而点赞通知可以接受 30 秒的延迟。更棘手的是连接管理千万 DAU 意味着同时维护百万级的 WebSocket 长连接连接断开、重连、App 切后台、设备网络切换——这些行为导致的连接状态变化必须在系统层面可靠处理否则消息丢失率会直线上升。本文复盘一套支持千万级设备的消息推送架构涵盖长连接管理、在线/离线分发策略、APNs/FCM 通道管理和推送到达率监控。二、整体推送架构长连接网关设计3.1 连接管理长连接网关使用 Netty 实现每个网关节点维护 5~10 万条 WebSocket 连接。核心组件Component public class WebSocketGateway { // 本节点维护的连接channelId → Channel private final ConcurrentHashMapString, 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; MapObject, 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 通道管理离线用户通过 APNsiOS或 FCMAndroid推送。通道管理的核心关注点是证书/密钥轮换和到达率监控Service public class OfflinePushService { private final MapString, ApnsClient apnsClients new ConcurrentHashMap(); private final MapString, 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 { PushNotificationResponseSimpleApnsPushBuilder 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. 获取用户所有设备的在线状态 ListDeviceInfo devices userDeviceService.getUserDevices( message.getUserId()); ListDeviceInfo onlineDevices new ArrayList(); ListDeviceInfo 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 在大陆的降级方案利用机器学习预测用户的最佳推送时机提高点击率以及构建推送策略引擎——根据消息类型、用户活跃度和时段动态选择推送通道和频率。

相关新闻