航班通知系统高并发推送架构设计:从事件到千万用户触达
“三千万用户”和“1.8 秒”这两个数字放在一起第一眼容易让人觉得是一个营销向的数据包装。但如果实际做过高并发通知系统的开发就会知道这是一个非常典型的工程问题不是“发一条消息”那么简单而是在航班延误、取消、登机口变更这种突发场景下把一个“事件”扩散成几百万甚至几千万条“用户通知”还要保证延迟可控。这个问题的难点不在某一台服务器而在整体链路的分层设计、订阅匹配、扇出批处理、第三方推送通道的适配以及极端情况下的兜底策略。这篇文章会围绕航班通知系统的架构设计展开讲清楚一条航变通知从产生到触达用户手机上中间会经过哪些环节也会给出可以落地的代码框架、配置示例和排错思路。不管你是正在做 OTA、航空、出行类业务的技术人还是想从系统设计角度理解高并发推送架构都能从里面拿到一些可复用的判断。先给一个明确结论1.8 秒绝不是某个单点功能做得好就能实现的数字它需要把整个通知链路拆成可量化的延迟预算然后对每一段进行优化。下面我们一步一步拆。1. 这件事的难点在哪里不是“发消息”而是“爆发式扇出”很多人第一次看到“三千万用户收到航班通知”时下意识会觉得这是推送平台的能力问题。只要接入友盟、极光、个推或者直接对接厂商通道把消息群发出去就可以了。但实际到了工程层面情况要复杂得多。首先航班通知不是无差别群发。同一个航班可能关联几千名乘客但这些乘客乘坐的是不同日期、不同舱位、不同渠道购买的机票他们对通知的需求也不同。有的用户关注延误有的用户关注登机口变更还有的用户只希望收到取消通知不希望被低优先级消息打扰。所以通知系统必须先做“订阅匹配”也就是把同一个航班事件匹配到真正需要通知的用户集合上。其次航班事件具有明显的突发性。正常情况下航班状态稳定系统压力很小。一旦遇到雷雨天气、空域管制或其他异常情况可能同时在短时间内有多个航班发生延误或取消受影响用户量会瞬间放大。这种“平时很闲峰值很高”的流量特征要求系统必须具备足够的峰值容量而不是按平均请求量去设计。第三通知内容是个性化的。同样一个航班延误A 用户的值机状态、B 用户是否已经托运行李、C 用户是否购买了退改签权益都会影响通知文案和服务建议。如果只是把同一句话复制给所有人产品体验会很差。所以扇出阶段不仅要找到用户还要找到该用户与本次事件相关的上下文。最后也是最容易被忽略的一点消息通道本身的限制。iOS 有 APNs安卓各个厂商有各自的推送服务国内第三方推送平台也各有能力边界。每个通道的批量发送接口、限流策略、回调机制、失败语义都不一样。通知系统必须在这一层做适配和降级才能保证整体延迟可控。所以可以用一句话总结航班通知系统的核心难点是把“一条事件”高效地转换为“多条个性化推送任务”并在多通道、多供应商、多失败场景下保证可靠性。三千万用户是一个规模数字但真正考验系统的是爆发式扇出能力。2. 一条航班通知的后台旅程从航变事件到用户手机为了让后面的架构讨论有共同语言我们需要先建立一张完整的链路图。实际上一条航班通知从产生到最终送达会经过以下环节。2.1 航变事件产生航班状态的变化并不由通知系统自己决定而是来自上游系统比如 AOC 运行控制中心、机场数据源、空管信息、航司内部航班管理系统等。这些系统会输出航班变更事件通常包括航班号、航班日期、事件类型、变更后的时间、登机口、航站楼信息等。这一层的关键是事件质量和实时性。如果上游数据延迟很大下游再优化也没有用。更麻烦的是上游系统可能重复推送同一个事件也可能推送的事件顺序混乱比如先收到“延误”再收到“恢复正点”最后又收到一次“新的延误”。2.2 事件标准化与去重通知系统接收到上游事件后首先需要做标准化也就是把不同来源、不同格式的数据统一成内部事件模型。这个阶段要处理两个问题一是去重。同一个事件上游可能推送了三次系统必须保证只生成一次通知任务。去重维度一般是“航班号航班日期事件类型数据源”有些场景还要加上事件版本号。二是防乱序。航变事件之间存在先后关系系统不能因为某条消息晚到就覆盖掉已经通知过用户的新状态。2.3 用户订阅匹配标准化之后系统需要回答一个关键问题哪些用户关心这个航班。通常业务库里有预订表、订单表、行程单表但这些表并不适合在推送链路的高频路径上去查询。更合理的方式是维护一份独立的“订阅关系缓存”提前把“航班号日期 - 用户 ID 集合”的关系建立好。这样一来航班事件到达时系统直接从缓存里取出用户列表而不是在 MySQL 里执行复杂 join 查询。基于 Redis 的 Set 结构是比较常见的做法也可以用内存缓存但要考虑多实例的一致性问题。2.4 生成通知任务并扇出拿到用户列表之后系统会为每一个需要通知的用户生成一条通知任务。任务里不仅包含用户 ID 和推送内容还要包含本次事件的 eventId、taskId、推送目标 token、优先级、可重试次数等信息。这一步会引入消息中间件。事件消费者把批量的用户 ID 拆分成多个任务或者按一定维度打包成批次写入 MQ。下游的推送 Worker 从 MQ 消费这些任务真正完成推到渠道的动作。使用 MQ 的好处是削峰填谷、失败重试、消费隔离。2.5 渠道适配与批量发送推送 Worker 消费到任务后会根据用户设备类型、推送通道、厂商限制等因素把任务分组然后调用对应的适配器发送。适配器要屏蔽不同通道的协议差异。这里有个重要的优化点批量。千万级任务如果一条一条调 HTTP 接口延迟和系统开销都不可接受。业界常见做法是每次批量发送 100 到 1000 条具体取决于第三方通道的限制。批量发送不仅能提升吞吐还能减少与上游建立连接的次数。2.6 回执、重试与补偿推送发送之后并不代表结束。系统需要处理回调结果哪些成功、哪些失败、哪些 token 已经失效。失败的任务要分情况处理可重试的进入重试队列不可重试的直接记录原因。比如用户卸载了 APPtoken 失效那么再重试也没有意义而第三方网关超时则多属于可重试场景。2.7 落库与统计分析最后一步是记录通知的发送状态、用户打开情况、点击情况。这些数据一方面用于业务分析另一方面也是监控告警的基础。只有知道每分钟推送量、成功率、平均延迟才能判断系统是否真正达标。从事件产生到用户手机收到通知整个链路的合理延迟预算大约可以这样分配环节延迟预算事件接入与标准化100ms 以内订阅匹配200ms 以内扇出与入队100ms 以内推送 Worker 消费与批处理300ms 以内第三方通道下发1s 左右用户手机展示100ms 以内如果题目中说的 1.8 秒是指后台链路从事件接收到推送网关完成下发那么上面这个预算表大致成立。如果是指端到端那还要加上第三方通道和手机厂商推送服务的处理时间通常会超过 1.8 秒。这也说明聊这一类指标前对齐统计口径比数字本身更重要。3. 系统设计的总原则三层解耦与四个约束看完链路之后我们再把视角抬高一点看整体架构应该怎么设计。一个能支撑千万级用户的航班通知系统一般会分为三层事件接入层、业务核心层、渠道适配层。事件接入层负责接收上游航班变更数据屏蔽数据源的差异输出标准化的内部事件。业务核心层负责订阅匹配、任务生成、扇出重试是整个系统的大脑。渠道适配层负责对接所有推送通道统一调用接口和回调处理。在具体落地时需要坚持四个约束。第一个约束是数据库不能成为热点路径的瓶颈。订阅匹配、token 查询、任务状态更新的高频读写都不应直接打到底层业务数据库。要么使用 Redis 做订阅缓存要么把热数据提前加载到内存或归档表中。否则一个航班事件过来系统会发现要同时几千次甚至几万次查询数据库很容易把连接池打爆。第二个约束是同步链路尽量短。航班事件处理链路里除了必要的 RPC 调用之外尽量不要有大量同步等待。推送发送尤其要注意不能在主线程里同步等待第三方通道返回否则吞吐会非常低。正确的做法是异步化接收事件时快速返回后续依赖 MQ 逐步推进。第三个约束是渠道适配必须独立。每个推送通道的协议不同、限流不同、错误码不同。如果渠道适配逻辑与业务核心逻辑耦合在一起后续每接入一个新通道都会改动核心代码风险很大。推荐的模式是定义一个统一的 PushGatewayClient 接口每次接入新通道就新增一个实现类。第四个约束是每条通知任务必须幂等。推送是天然存在重复风险的场景MQ 可能重复消费第三方通道可能超时重试网络重传也可能导致重复。如果系统不做幂等用户可能同一个航班收到四五条一模一样的延误通知产品体验会非常差。最常见的做法是使用唯一的 taskId 作为幂等键在 Redis 里做重复判断。4. 技术选型消息中间件、缓存、推送通道该如何权衡关于具体技术栈很多团队会有不同选择但在航班通知这个业务里核心角色是确定的消息中间件、缓存、推送通道、配置中心、监控系统。下面分别说清楚选型思路。4.1 消息中间件消息中间件在系统里承担两个职责一个是接收上游航班变更事件另一个是承载扇出后的推送任务。事件接入场景对顺序要求不高但要求吞吐量大、支持重复消费排查。Kafka 和 RocketMQ 都是不错的选择。RocketMQ 对事务消息、顺序消息、延迟消息有更完整的支持在航班取消后的延迟补救任务中比较方便Kafka 吞吐更高生态更成熟团队熟悉度通常也更好。扇出任务的场景更关注消费堆积能力和失败重试能力。这里强烈建议不要把所有任务都放到一个 Topic 里而是按航线、航班日期或用户分片拆成多个 Topic 或者多个队列。这样可以避免一个消费组整体被慢任务拖住。选择建议如果团队已经熟悉 Kafka事件接入和扇出任务都可以使用 Kafka通过配置不同的 topic 和消费组隔离。如果追求更细的控制能力可以用 RocketMQ它在告警、延迟消息、事务消息方面能省不少研发成本。4.2 缓存缓存主要做两件事存储订阅关系和记录幂等键。订阅关系的特点是读多写少适合放在 Redis 的 Set 结构中一个 key 对应一个航班的所有用户。幂等键的特点是短期有效使用 setnx 命令即可。需要特别注意的是订阅关系缓存的更新不能只靠业务下单时写入。因为用户可能取消订单、改签、值机后更换手机号缓存和业务库之间的一致性需要靠消息订阅或者定时任务来维护。否则很可能出现用户已经退票却还能收到航班通知的情况。4.3 推送通道国内安卓推送的现状是系统级推送能力分散在华为、小米、OPPO、vivo 等厂商手里第三方聚合平台可以通过厂商通道提高送达率。iOS 则统一走 APNs。一个完整的航班通知系统通常会在第三方聚合平台之上再叠加厂商直连通道作为高优先级通知的专用路径。选择通道时要重点考察三个指标送达率、必达能力、限流策略。送达率决定用户体验必达能力决定航班取消这类高优先级事件能否在系统休眠时也弹出通知限流策略决定系统峰值时是否会被上游掐断。4.4 配置中心和监控配置中心负责管理推送开关、批量大小、重试次数、黑白名单等动态配置。航班大规模延误时运营可能需要临时关闭非紧急通知或者下调某一通道的发送速率。如果这些参数写死在代码里发布一次可能要好几分钟这在线上是不可接受的。监控系统则需要覆盖从事件进入到推送回执的全链路。最少要监控事件接入延迟、订阅匹配耗时、MQ 消费 Lag、推送成功率、通道调用平均耗时、P95 延迟。没有监控1.8 秒只能是一个没有依据的口号。5. 核心流程拆解从事件到达开始在动手写代码之前先把核心流程完整拆一遍。这个流程会直接对应后面代码示例中的模块划分。5.1 步骤一接收并校验航班事件事件消费者从 MQ 中读取航班变更消息首先做格式校验。字段不完整、事件时间晚于当前时间过久、事件来源不在白名单内这些都算非法事件直接丢弃并记录告警。合法事件进入下一步。这一步是防御性的看起来简单但在真实生产环境中非常重要。上游系统并不总是值得信任有些时候一个坏数据会引发整个通知链路的异常。5.2 步骤二根据航班号匹配订阅用户使用航班号和航班日期作为 key从 Redis 中取出订阅了该航班的用户 ID 列表。这一步有两种实现一种是把用户 ID 全量取出直接交给下游。优点是简单直接缺点是当用户量达到几万甚至几十万时列表很大传输和写入 MQ 都有压力。另一种是分段处理比如每 1000 个用户拆成一个批次分批写入 MQ。这样的好处是可以控制消息体大小降低消费者压力也更利于并行处理。推荐使用第二种方式。后面代码示例也会采用分批写入的逻辑。5.3 步骤三构建推送任务对每一个用户需要构建一条推送任务。任务里通常包含以下字段字段说明taskId全局唯一用于幂等判断userId用户 IDeventId对应的航班事件 IDflightNo航班号flightDate航班日期eventType事件类型如延误、取消、登机口变更pushToken用户设备的推送 tokenchannelType推送通道类型priority优先级决定是否插队这里要注意pushToken 不应该在扇出阶段临时去查数据库。更合适的做法是在用户注册订阅时就把 token 一起写入订阅缓存。否则扇出阶段会因为查询用户 token 而增加大量延迟。5.4 步骤四推送 Worker 消费任务并进行批量发送推送 Worker 从 MQ 中拉取任务按照通道类型和批次大小进行聚合。比如同一个通道的 500 个任务合并成一个 batch调用一次 PushGatewayClient 的批量发送接口。这个环节常见的性能瓶颈有两个。一个是一条任务一个线程导致线程数膨胀CPU 大量消耗在线程切换上。另一个是没有限制并发当上游事件爆发时下游第三方通道被瞬间打满触发限流。更合理的做法是使用有界队列和固定线程池同时根据第三方通道的返回状态动态调整发送速率让系统具备一定的自适应能力。5.5 步骤五处理发送结果与重试批量发送完成后系统会得到一个结果对象。对于成功任务更新发送状态做数据落库。对于失败任务判断是否可重试。可重试任务进入延迟重试队列超过重试限次的进行告警。重试策略要避免无脑重试。比如第三方通道返回“应用被卸载”“token 无效”这是不可重试错误应该直接标记失败。只有超时、限流、服务端异常等情况才适合重试。5.6 步骤六兜底通道有一些场景下APP 推送并不能保证触达用户。比如用户关闭了通知权限、手机处于离线状态、APP 长时间未打开。此时需要短信、邮件甚至电话语音等兜底通道。但兜底通道的成本较高且有被运营商限制的风险。一般只会选择高优先级事件使用兜底通道比如航班取消、航班合并等。普通登机口变更只通过 APP 推送即可。6. 可落地的核心代码与最小验证前面讲了很多设计原则这一节给出一个可以跑通最小链路的代码框架。考虑到不同团队技术栈不一样下面的示例以 Spring Boot Kafka Redis 为例换成 RocketMQ 或 RabbitMQ 后思路仍然一致。6.1 项目依赖与配置如果是 Maven 项目最小依赖如下dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency /dependencies对应的 application.yml 配置如下# 文件路径src/main/resources/application.yml spring: application: name: flight-notify redis: host: ${REDIS_HOST:127.0.0.1} port: ${REDIS_PORT:6379} kafka: bootstrap-servers: ${KAFKA_SERVERS:127.0.0.1:9092} consumer: group-id: notify-center enable-auto-commit: false auto-offset-reset: latest listener: ack-mode: manual notify: fanout: batch-size: 500 push: worker-count: 8 max-retry: 3这里把自动提交关掉改为手动 ack可以减少消费成功但未提交导致的重复推送问题。监听器 ack-mode 设置为 manual是为了在业务处理完成后手动确认偏移量。6.2 航班事件消费者下面这个类消费上游的航班事件。核心逻辑是校验事件、获取订阅用户、分批生成扇出任务。// 文件路径src/main/java/com/travel/notify/consumer/FlightEventConsumer.java Slf4j Component public class FlightEventConsumer { private final FlightSubscriptionService subscriptionService; private final KafkaTemplateString, Object kafkaTemplate; private final int batchSize; public FlightEventConsumer(FlightSubscriptionService subscriptionService, KafkaTemplateString, Object kafkaTemplate, Value(${notify.fanout.batch-size:500}) int batchSize) { this.subscriptionService subscriptionService; this.kafkaTemplate kafkaTemplate; this.batchSize batchSize; } KafkaListener(topics flight-event, groupId notify-center) public void onFlightEvent(FlightChangedEvent event, Acknowledgment ack) { if (event null || !StringUtils.hasText(event.getFlightNo())) { log.warn(invalid flight event); ack.acknowledge(); return; } // 根据航班号和日期获取已订阅用户 ListString userIds subscriptionService.findSubscribeUserIds( event.getFlightNo(), event.getFlightDate()); if (userIds.isEmpty()) { ack.acknowledge(); return; } // 按固定大小拆批避免单个消息过大 for (int i 0; i userIds.size(); i batchSize) { int end Math.min(i batchSize, userIds.size()); ListString subList userIds.subList(i, end); PushFanoutTask task PushFanoutTask.builder() .taskId(UUID.randomUUID().toString().replace(-, )) .eventId(event.getEventId()) .flightNo(event.getFlightNo()) .flightDate(event.getFlightDate()) .eventType(event.getEventType()) .userIds(subList) .build(); kafkaTemplate.send(push-fanout-task, task.getFlightNo(), task); } ack.acknowledge(); log.info(flight event handled, flightNo{}, userIds{}, event.getFlightNo(), userIds.size()); } }这段代码有一个容易踩坑的地方把 userIds 列表直接放进消息体。如果航班量级很大比如一个航班关联几万用户那么拆批后的消息体仍然会很大。建议在实际场景中对 userIds 做压缩或者只传一个“用户批次查询条件”下游再去查询完整的用户信息。为了把示例控制在可读范围这里保留了列表方式。6.3 订阅关系查询订阅关系查询依赖 Redis。这里实现一个简单的读取逻辑以 flightNo 和 flightDate 组成 key在 Set 中取出用户集合。// 文件路径src/main/java/com/travel/notify/service/FlightSubscriptionService.java Slf4j Service public class FlightSubscriptionService { private final StringRedisTemplate stringRedisTemplate; public FlightSubscriptionService(StringRedisTemplate stringRedisTemplate) { this.stringRedisTemplate stringRedisTemplate; } public ListString findSubscribeUserIds(String flightNo, String flightDate) { String key sub:flight: flightNo : flightDate; SetString members stringRedisTemplate.opsForSet().members(key); if (members null || members.isEmpty()) { return List.of(); } return new ArrayList(members); } }需要注意OpsForSet 的 members 方法会一次性取出整个 Set。当订阅用户量特别大时建议改用 sScan 游标遍历避免 Redis 阻塞。在实际项目中还需要通过消息订阅维护这个 Set比如用户下新单时加入退票时移除。6.4 推送 Worker 与幂等去重推送 Worker 从push-fanout-task这个 topic 中消费扇出任务调用渠道适配层发送。核心有两个批量发送和幂等去重。// 文件路径src/main/java/com/travel/notify/worker/PushWorker.java Slf4j Component public class PushWorker { private final PushGatewayClient pushGatewayClient; private final DedupService dedupService; public PushWorker(PushGatewayClient pushGatewayClient, DedupService dedupService) { this.pushGatewayClient pushGatewayClient; this.dedupService dedupService; } KafkaListener(topics push-fanout-task, groupId notify-push-worker) public void onFanoutTask(PushFanoutTask task) { String taskId task.getTaskId(); // 幂等判断如果已经处理过直接跳过 if (!dedupService.isFirstTime(taskId)) { log.info(duplicate task, skip taskId{}, taskId); return; } ListPushTask pushTasks buildPushTasks(task); PushBatchResult result pushGatewayClient.batchSend(pushTasks); if (result.hasFailItems()) { for (PushTask failedTask : result.getFailedTasks()) { if (failedTask.isRetryable()) { // 实际项目中可写入延迟重试 topic log.warn(retryable task failed, taskId{}, reason{}, failedTask.getTaskId(), failedTask.getFailReason()); } else { // 不可重试记录失败原因 log.warn(unretryable task failed, taskId{}, reason{}, failedTask.getTaskId(), failedTask.getFailReason()); } } } } private ListPushTask buildPushTasks(PushFanoutTask task) { // 根据 token 明细构造具体推送任务此处为示例 return task.getUserIds().stream() .map(userId - PushTask.builder() .userId(userId) .flightNo(task.getFlightNo()) .eventType(task.getEventType()) .taskId(task.getTaskId() : userId) .build()) .collect(Collectors.toList()); } }幂等去重的实现依赖 Redis 的 setnx 操作// 文件路径src/main/java/com/travel/notify/service/DedupService.java Slf4j Service public class DedupService { private final StringRedisTemplate stringRedisTemplate; public DedupService(StringRedisTemplate stringRedisTemplate) { this.stringRedisTemplate stringRedisTemplate; } public boolean isFirstTime(String taskId) { String key notify:dedup: taskId; return Boolean.TRUE.equals( stringRedisTemplate.opsForValue().setIfAbsent(key, 1, Duration.ofHours(24)) ); } }这个 setnx 方案的关键点有两个一是任务 ID 必须全局唯一二是过期时间要合理。过期时间太短重复消费可能会发生在过期之后过期时间太长Redis 内存压力会增大。通常设置为 24 小时即可覆盖大多数重复消费场景。6.5 最小链路验证方式在本地启动 Kafka 和 Redis 后可以先用一个简化的 Web 接口手动触发事件模拟一次航班变更。curl -X POST http://127.0.0.1:8080/events/simulate \ -H Content-Type: application/json \ -d { eventId: E20250101001, flightNo: CA1234, flightDate: 2025-01-01, eventType: DELAY, newDepartureTime: 2025-01-01T22:30:00 }这里的模拟接口只是演示用实际生产环境并不建议开放。更常见的做法是通过 Kafka 控制台生产命令投递消息kafka-console-producer.sh --broker-list 127.0.0.1:9092 --topic flight-event {eventId:E20250101001,flightNo:CA1234,flightDate:2025-01-01,eventType:DELAY,newDepartureTime:2025-01-01T22:30:00}消息投递成功后观察应用日志。正常情况下可以依次看到航班事件消费者打印“flight event handled”并统计 userIds 数量。推送 Worker 打印“duplicate task”或成功发送的日志。如果上游故意重复投递同一条消息第二次会被幂等逻辑拦截不会重复推送。这个最小链路虽然简单但已经覆盖了事件接入、订阅匹配、扇出、幂等、批量发送这些核心机制。后续要做线上验证只需要替换为真实的推送通道适配器并补充监控埋点。7. “1.8 秒”如何验证SLO 与监控设计代码写完了怎么证明系统能达到 1.8 秒的延迟目标答案不是凭空估算而是通过监控数据验证。在讨论这个问题之前必须先明确一个口径1.8 秒到底是“后台链路从事件接入到推送网关完成”的耗时还是“事件发生到用户实际看到通知”的端到端耗时。如果是后台链路目标那么关键指标可以定义为event_receive_time上游事件进入通知系统的时间。push_send_time推送 Worker 调用第三方通道接口完成发送的时间。二者差值即为后台链路延迟。如果是端到端目标那还需要加上第三方通道的处理时间以及 App 进程收到厂商离线消息后展示通知的时间。这部分往往不是通知系统能完全控制的它依赖用户手机状态和厂商推送服务的时效。从系统设计的角度我更建议把 SLO 拆成三段来管理指标目标事件接入延迟 P95200ms订阅匹配与扇出延迟 P95300ms推送 Worker 批量发送延迟 P95800ms端到端延迟 P95与各通道质量相关监控埋点可以在代码里直接实现。最简单的办法是给每个任务增加一个 eventCreateTime 字段从事件产生时就把时间戳透传到推送任务中。推送发送完成时计算当前时间与事件产生时间的差值。// 伪代码计算通知链路延迟 long now System.currentTimeMillis(); long eventCreateTime task.getEventCreateTime(); long backendLatency now - eventCreateTime; metrics.timer(notify.backend.latency, eventType, task.getEventType()) .record(backendLatency, TimeUnit.MILLISECONDS);有了这个埋点就能在 Grafana 里看到从事件接入到推送完成的分位分布。如果 P95 超过 1.8 秒排查思路可以按照链路逐段查看是事件消费者消费太慢还是上游投递本身有延迟是订阅查询 Redis 耗时高还是任务入队积压是推送 Worker 消费能力不足还是第三方通道限流在实际项目中经验教训是延迟问题很少出在单段代码逻辑上更多是出现在 MQ 积压、线程池阻塞、第三方通道排队这三类问题上。所以监控面板上需要重点展示各环节 MQ 的消费 Lag、线程池队列长度、以及第三方通道的耗时分布。8. 常见问题与排查方法即使架构设计得再完善生产环境还是会出现各种问题。这里整理几个典型的故障场景和排查思路。问题现象可能原因排查方式解决方案高峰期 MQ 大量积压下游消费者处理能力不足查看消费 Lag、消费者线程数、通道平均耗时增加消费实例批量大小调大或动态降低非核心通道发送速率用户收到重复通知MQ 重复消费或上游重复推送检查幂等日志确认 Redis 去重键是否生效校验幂等键设计和过期时间统一 taskId 生成规则通知延迟突然升高第三方通道限流或 token 失效导致重试过多查看通道调用失败率、重试链路耗时对通道做熔断和降级优先保证高优先级事件某个航班事件触发后用户量异常大订阅缓存脏数据包含退票或改签用户对照订单系统检查缓存写入逻辑增加缓存更新消息定时做全量对账Worker 线程被阻塞同步调用第三方通道且未设置超时检查线程池 dump 和 HTTP 连接池配置统一设置连接超时和读取超时使用异步 Http 客户端推送成功率低但监控无告警监控只统计发送成功没有统计到达成功接入推送通道回执数据建立成功率、到达率双层监控设置告警阈值这里特别想强调一个容易被忽略的问题Android 推送的“到达率”和“成功率”不是一个概念。第三方通道返回“发送成功”并不代表用户真正看到了通知有些手机厂商会合并通知、折叠通知甚至因为省电策略杀掉 App 进程。如果产品指标是“用户看到通知”那么监控就必须以厂商回执的“到达”为准而不是以网关返回的成功为准。另一个常见问题是通知内容里的时间格式。航班延误后新的出发时间必须按用户所在时区或航班起飞机场所在时区进行格式化。如果系统只存了一个 UTC 时间没有存时区信息用户看到的时间就可能是错的。这类问题不常发生但一旦发生用户投诉会比较集中。9. 生产环境最佳实践与最后的取舍9.1 高优先级事件与普通事件分级航班取消这种通知必须做到“快”和“必达”。而登机口变更、延误提醒这类通知则属于普通通知允许一定延迟。实现分级的方法有很多使用不同优先级的 MQ topic或者在任务字段中标记优先级推送 Worker 消费时优先处理高优先级任务。项目落地时我强烈建议从第一天就用独立 topic 隔离这两种任务。否则一旦遇到大面积航班取消普通任务会占据消费能力导致真正紧急的通知无法及时送达。分级还有一个好处高优先级事件可以把兜底通道短信、电话语音打开普通事件不做兜底成本可控。9.2 弹性扩缩容航班通知系统属于典型的“峰值突发型”业务可以充分利用弹性扩缩容。预警机制可以基于 MQ 消费 Lag 和 Redis 中待处理任务量来判断。当消费 Lag 超过阈值时自动增加推送 Worker 实例流量回落后再收缩。这里要注意的是消费者实例不是越多越好。推送 Worker 的瓶颈很可能不在自身计算能力而在于第三方通道的每分钟请求配额。如果实例数增加导致并发请求翻倍第三方通道会直接限流效果适得其反。更稳妥的做法是增加实例的同时在代码里动态调节每个实例的发送速率让整体 QPS 稳定在通道允许范围内。9.3 数据安全与合规航班通知系统涉及大量用户行程数据、手机号、设备 token 和航班信息数据安全必须做到位。第一用户设备 token 是非常敏感的信息存储时必须加密。它本身可以作为推送的唯一凭据一旦泄露攻击者可以向用户发送恶意推送。数据库中的 token 字段不能明文保存日志中更不能打印。第二订阅关系缓存里尽量不要放手机号等明感信息。匹配用户时只需要 userId 和 token手机号只在短信兜底环节使用且需要走独立的权限审批流程。查询、导出、接口调用都应该有审计日志。第三推送内容本身也要防止泄露。航班通知中会包含旅客姓名、行程、航班号这些都属于个人出行数据。调用第三方推送通道时建议对内容做加密或使用通道提供的私密消息能力避免明文经过第三方侧。尤其在对接非官方通道时这一点需要格外关注。9.4 故障演练通知系统是典型的“平时看不出问题出事才见真章”的系统。建议定期组织故障演练重点模拟三种场景第一种是大面积航班取消。可以提前准备一批模拟航班事件以正常峰值的三倍量压入系统观察 MQ 是否积压、第三方通道是否限流、短信通道是否被冻结。第二种是第三方通道长时间不可用。切断某个主要推送通道看系统能否自动切换备用通道用户通知的到达率会下降到什么程度。第三种是 Redis 故障或重建。订阅关系和幂等键都依赖 Redis如果 Redis 不可用系统是否具备降级方案。比如订阅关系可以临时切换到本地缓存或数据库兜底幂等判断可以暂时放行但记录标记等待 Redis 恢复后做一次去重清理。这些演练看起来耗时但价值很大。很多系统就是在故障演练中才发现架构里的单点问题和配置错误。如果没有演练真正发生故障时团队只能摸着石头过河延迟、成功率、用户投诉都会失控。9.5 主要的取舍如果让我总结这类系统里最重要的几个取舍我会记住这样几条先定义清楚“1.8 秒”的含义再开始做性能优化。没有统一口径的指标是没有意义的。用 MQ 做削峰填谷是对的但一定要监控消费 Lag。延迟通常不是代码慢而是积压。订阅关系缓存是扇出性能的关键但缓存和业务库的一致性必须通过异步对账持续保证。推送通道越多系统稳定性越强但复杂度也越高。每个通道背后都要有配额、限流、熔断和降级策略。幂等是推送系统的基本功。宁可在前面多花精力设计 taskId 和去重方案也不要在线上被重复推送投诉打爆。航班通知系统看似只是一个“发消息”的服务但要做到千万用户规模下稳定快速触达需要把事件接入、订阅匹配、扇出、批量发送、重试、监控和兜底串成一条完整的链路。只要能把这套链路跑通这套系统设计方法也可以迁移到外卖、电商、物流、出行等几乎所有需要个性化通知的业务中。