尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
视频平台的消息推送架构:从长连接到离线推送的高可用方案
视频平台的消息推送架构从长连接到离线推送的高可用方案一、背景与问题定义视频平台的消息推送场景远比即时通讯复杂。用户可能收到互动通知评论、点赞、关注、系统通知审核结果、活动推送、以及实时消息直播开播提醒。这些场景对时效性和可靠性的要求各不相同——直播开播提醒需要在 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 在大陆的降级方案利用机器学习预测用户的最佳推送时机提高点击率以及构建推送策略引擎——根据消息类型、用户活跃度和时段动态选择推送通道和频率。
RELATED

相关推荐

少儿编程入门:从C++基础语法到第一个交互程序实战

少儿编程入门:从C++基础语法到第一个交互程序实战

1. 项目概述:为什么从C开始?如果你正在为孩子寻找一门编程入门语言,或者你自己就是一位希望引导孩子进入编程世界的家长、老师,面对Python、Scratch、C这些选项,可能会有些犹豫。Scratch图形化,上手快&…

📅 2026/7/26 1:18:12
网关层的安全防护体系建设——从 WAF 到 API 安全的纵深防御架构

网关层的安全防护体系建设——从 WAF 到 API 安全的纵深防御架构

网关层的安全防护体系建设——从 WAF 到 API 安全的纵深防御架构 一、网关安全不能只靠一层 WAF 很多团队对网关安全的理解是"在 Nginx 前面加一个 WAF 就差不多了"。WAF 确实能拦截 SQL 注入、XSS、CSRF 等常见的 Web 攻击,但现代 API 的安全威胁已经远远…

📅 2026/7/26 9:50:41
JavaSE基础概念笔记01

JavaSE基础概念笔记01

数值取值范围从小到大 byte<short<int<long<float<double 隐式转换(自动类型提升):取值范围小的数据转换成取值范围大的数据 记忆:偷偷变强(即为隐式) 1.取值范围小的数据和取值范围大的数据进行运算时,小的会先转换为大的,再进行运算 2.byte char short 进行运…

📅 2026/7/27 18:01:39
MORE NEWS

更多资讯

📰

AI时代PLC工程师的生存法则:从写代码到搞定产线

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

MDPI投稿状态全解析:11个状态含义、时间线与催稿技巧

1. 投稿状态到底在说什么第一次往MDPI旗下期刊投论文的人&#xff0c;十有八九会被投稿系统里那一串状态搞得心里七上八下。Submitted、Under Review、Pending Decision、Accepted……每个词都认识&#xff0c;但连在一起就不知道到底进展到哪一步了。更让人焦虑的是&#xff0…

📰

SpringBoot+Android民宿预订系统从零到答辩全指南

简介&#xff1a;一份基于Spring Boot与Android平台的民宿预订系统毕业论文文档&#xff0c;面向计算机相关专业毕业生以及需要完成课程设计或毕业设计的开发者。内容系统阐述了民宿预订系统的设计目的、需求分析、总体架构与实现方案&#xff0c;重点涉及Spring Boot框架选型、…

📰

ValidX校验库集成指南:Maven/Gradle构建与镜像配置

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

RocksDB 文档站深度指南:docs 目录 Jekyll 站点的结构、配置与定制方法

RocksDB 文档站深度指南&#xff1a;docs 目录 Jekyll 站点的结构、配置与定制方法 【免费下载链接】rocksdb A library that provides an embeddable, persistent key-value store for fast storage. 项目地址: https://gitcode.com/gh_mirrors/ro/rocksdb 本文围绕 Ro…

📰

ZeroClaw Agent 确定性回放评测:用 zeroclaw-eval 验证 Agent 机制的正确性

ZeroClaw Agent 确定性回放评测&#xff1a;用 zeroclaw-eval 验证 Agent 机制的正确性 【免费下载链接】zeroclaw Fast, small, and fully autonomous AI personal assistant infrastructure, any OS, any platform — deploy anywhere, swap anything &#x1f980; 项目地…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬