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

相关新闻

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

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

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

2026/7/23 8:19:11 阅读更多 →
AI技术如何重塑现代教育:学生主导的变革

AI技术如何重塑现代教育:学生主导的变革

1. 教育变革的新浪潮:当学生成为技术革命的主角 三年前我在某高校做技术分享时,台下学生还在用纸笔记录;去年再去,满眼都是学生用AI工具实时转录、思维导图同步的场景。这种肉眼可见的变化正在全球校园里加速发生——不同于以往由…

2026/7/23 8:18:10 阅读更多 →
TM4C1299NCZAD CAN控制器原理与实战:从帧结构到消息对象配置

TM4C1299NCZAD CAN控制器原理与实战:从帧结构到消息对象配置

1. CAN控制器核心原理与帧结构深度解析控制器局域网,也就是我们常说的CAN总线,在汽车电子和工业控制领域几乎是“基础设施”一样的存在。它不像我们熟悉的UART或I2C那样简单直接,其设计哲学从一开始就瞄准了高可靠性、实时性和多节点竞争的场…

2026/7/23 8:18:10 阅读更多 →

最新新闻

月嫂服务中心低成本获客神器,凡科全新1折优惠渠道:99做小程序只认餐宝盈,含零代码SAAS、AI编程、源码定制交付

月嫂服务中心低成本获客神器,凡科全新1折优惠渠道:99做小程序只认餐宝盈,含零代码SAAS、AI编程、源码定制交付

月嫂服务怎么获客?凡科全新1折优惠渠道:99做小程序只认餐宝盈 摘要 对月嫂服务商家来说,获客难点往往不在于门店没有服务能力,而在于线上表达弱、承接入口散、用户看见后不容易直接成交。现在,餐宝盈官网 cby888.com…

2026/7/23 13:10:23 阅读更多 →
2026年经典爬虫案例专栏|第1篇:Python爬虫基础入门与环境搭建

2026年经典爬虫案例专栏|第1篇:Python爬虫基础入门与环境搭建

1.1 爬虫技术概述 1.1.1 什么是网络爬虫 网络爬虫(Web Crawler),也称为网页蜘蛛(Spider)、网络机器人(Bot),是一种自动化程序,用于按照一定的规则自动浏览万维网(WWW),并获取网页内容。爬虫技术在当今互联网时代扮演着至关重要的角色,它是搜索引擎、数据挖掘、内…

2026/7/23 13:10:23 阅读更多 →
无需高成本碳源!Next Materials:闪蒸焦耳热让水葫芦变身 Cr(III) 高效吸附剂

无需高成本碳源!Next Materials:闪蒸焦耳热让水葫芦变身 Cr(III) 高效吸附剂

1. 背景重金属污染具有毒性强、难降解、易累积等特点,是工业废水治理中的长期难题。Cr(III) 常见于皮革鞣制、电镀、钢铁和纺织等过程,传统处理方法往往面临成本高、操作复杂和二次污染等限制。吸附法因工艺简单、材料可设计性强而具有应用潜力&#xff…

2026/7/23 13:10:23 阅读更多 →
无需长时热处理!ACS Applied Energy Materials:0.5 秒焦耳热让碳纤维/MnOx 电极快速成型

无需长时热处理!ACS Applied Energy Materials:0.5 秒焦耳热让碳纤维/MnOx 电极快速成型

1. 背景超级电容器兼具高功率密度、快速充放电和长循环寿命,是连接传统电容器与电池的重要储能器件。其中,锰氧化物(MnOx)因理论电容高、资源丰富、成本低和环境友好而受到关注,但其本征电导率低、离子传输受限、结构稳…

2026/7/23 13:10:23 阅读更多 →
东华大学Nano Res. Energy:退役锂电池石墨负极回收正在从粗放重构走向精准再生

东华大学Nano Res. Energy:退役锂电池石墨负极回收正在从粗放重构走向精准再生

1. 背景随着新能源汽车和储能产业快速扩张,锂离子电池退役潮正在把石墨负极推向资源循环的关键位置。石墨不仅是锂电负极的主流材料,也是被视为供应链安全和低碳制造共同约束的关键矿物;天然石墨资源分布集中,人工石墨制备又依赖高…

2026/7/23 13:10:23 阅读更多 →
AI 电动绞肉机智能驱动 覆盖主电机驱动、刹车控制、智能逻辑管理的高效选型方案

AI 电动绞肉机智能驱动 覆盖主电机驱动、刹车控制、智能逻辑管理的高效选型方案

AI 智能绞肉机(带称重、变速、自动保护、智能启停)对驱动电机控制提出新挑战:高扭矩启停、变速平稳、低噪音、高效率。微碧半导体基于先进 Trench 工艺,为您提供覆盖主电机驱动、刹车控制、智能逻辑管理的完整 AI 绞肉机功率解决方…

2026/7/23 13:09:22 阅读更多 →

日新闻

从单点好评到指数级传播:AI副业主理人必须掌握的4层口碑渗透模型(含ROI测算表)

从单点好评到指数级传播:AI副业主理人必须掌握的4层口碑渗透模型(含ROI测算表)

更多请点击: https://intelliparadigm.com 第一章:从单点好评到指数级传播:AI副业主理人必须掌握的4层口碑渗透模型(含ROI测算表) 当AI副业主理人不再仅满足于单次服务交付,而是主动构建可复用、可裂变、可…

2026/7/23 0:00:25 阅读更多 →
AI写作开头钩子设计:为什么你的AI文案完读率不足18%?——基于2,346篇A/B测试报告的归因分析

AI写作开头钩子设计:为什么你的AI文案完读率不足18%?——基于2,346篇A/B测试报告的归因分析

更多请点击: https://codechina.net 第一章:AI写作开头钩子设计:为什么你的AI文案完读率不足18%?——基于2,346篇A/B测试报告的归因分析 在对2,346篇跨行业AI生成文案的A/B测试数据进行聚类分析后,我们发现&#xff1…

2026/7/23 0:01:26 阅读更多 →
Chitchatter完整指南:免费开源的终极点对点安全聊天工具

Chitchatter完整指南:免费开源的终极点对点安全聊天工具

Chitchatter完整指南:免费开源的终极点对点安全聊天工具 【免费下载链接】chitchatter Secure peer-to-peer chat that is serverless, decentralized, and ephemeral 项目地址: https://gitcode.com/gh_mirrors/ch/chitchatter Chitchatter是一款革命性的安…

2026/7/23 0:01:26 阅读更多 →

周新闻

Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中,我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源,还是配置文件、证书等,都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下,但这…

2026/7/22 8:58:19 阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP(轻量级目录访问协议)作为企业级身份认证的黄金标准,已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时,发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/22 19:43:43 阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击: https://intelliparadigm.com 第一章:AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”,而是以可解释、可审计、可迭代的方式,赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/22 12:54:44 阅读更多 →

月新闻