百度360必应搜狗淘宝本站头条
当前位置:网站首页 > 技术教程 > 正文

如何使用SSE实现服务器高性能实时发送通知

mhr18 2024-12-08 15:04 29 浏览 0 评论

技术方案

为了实现的实时通知系统,实现了以下内容:

  • Server-Sent Events (SSE) :使用 SSE 在服务器和客户端之间建立实时、单向的通道,允许在通知发生时立即传递。
  • Redis Pub/Sub: 使用 Redis pub/sub 作为消息系统,可以让多个发布者向多个订阅者发送消息。
  • Notification Write API: 开发了一个通知写入 API 来收集来自其他应用程序的通知, 此 API 向 Redis 发送通知。
  • Notification Read API: 也开发了通知读取API,订阅Redis,监听所有的通知。 收到通知后,读取 API 将它们发送到通过 SSE 连接到它的客户端。

为什么选择SSE

1.HTTP 轮询虽然是一个可行的选项,但不提供即时通知传递。 使用轮询,必须以特定时间间隔对服务器执行 ping 操作以检查是否有新事件,这可能会占用大量资源且效率低下,尤其是在许多用户订阅了通知的情况下。

2.另一方面,SSE 在服务器和客户端之间提供实时、单向的通道,允许在通知发生时立即传送,提供了实时性的保障。

3.虽然 WebSocket是实时通信的实现方案之一,但与 SSE 相比,它们需要更复杂的实现,虽然允许服务器和客户端之间进行双向通信,但由于只需要向客户端推送通知,因此选择 SSE 作为更简单、更高效的解决方案。

Redis 发布/订阅

为了从多个应用程序 pod 发送通知,使用了 Redis pub/sub。 Redis pub/sub 是一个消息系统,可以让多个发布者向多个订阅者发送消息。

import { Redis } from 'ioredis';

@Injectable()
export class RedisService implements OnModuleInit {

  ...

  async onModuleInit() {
    this.subscriber = new Redis({
      sentinels: this.redisSentinelConfig.addresses,
      name: this.redisSentinelConfig.masterName,
      password: this.redisSentinelConfig.password,
    });
    
    await this.subscriber.subscribe(this.redisSentinelConfig.channelName, async (err, count) => {
      if (err) {
        this.logger.error(`Failed to subscribe: ${err.message}`);
        return;
      }
      this.logger.log(`Subscribed successfully! This client is currently subscribed to ${count} channels`);
    });
    
    this.subscriber.on('message', async (channel, message: string) => {
      const liveNotification: LiveNotification = JSON.parse(message);
      await this.liveNotificationService.emit(liveNotification);
    });
  }
}

系统设计

为了收集来自其他应用程序的通知,开发了Notification Write API,此 API 向 Redis 发送通知。 Notification Read API订阅Redis,监听所有通知。 收到通知后,读取 API 将通过 SSE把消息发送到 连接的客户端。

由于多个读取 API 可能运行在不同的应用程序 Pod 中,因此每个读取 API 都会收到来自 Redis 的所有通知。Notification Read API将根据连接到当前 pod 的客户端过滤通知。 确保只向每个连接的客户端发送相关通知。


import { AuthGuard } from '@nestjs/passport';
import { Controller, Sse, UseGuards } from '@nestjs/common';

@Controller('/live-notification')
export class LiveNotificationController {
  constructor(private readonly liveNotificationService: LiveNotificationService) {}

  @Sse()
  @ApiBearerAuth()
  @UseGuards(AuthGuard('jwt'))
  public getEventsBySeller(@Tracers() tracers: ITracers, @SellerId() sellerId: number) {
    return this.liveNotificationService.subscribeForSeller(sellerId);
  }
}
 

import { EventEmitter } from 'events';
import { filter, fromEvent } from 'rxjs';

@Injectable()
export class LiveNotificationService implements OnModuleInit {
  
  private readonly emitter = new EventEmitter();
  
  ...
  
  public async emit(data: LiveNotification) {
    this.emitter.emit('liveNotification', { data });
  }
  
  public subscribeForSeller(sellerId: number) {
    const source = fromEvent(this.emitter, 'liveNotification');
    return source.pipe(
      filter(({ data: liveNotification }) => 
        liveNotification?.content == 'heartbeat' || 
        liveNotification?.sellerId == sellerId)
    );
  }
}
    

心跳消息

为避免服务器和浏览器之间的连接关闭,需要发送心跳消息来检测断开连接。为了解决这个问题,实现了每 30 秒发送一次心跳消息,从而解决了这个问题。

onModuleInit (): any { 
  setInterval ( () => { 
    const emitterListenerCount = this . emitter . listenerCount ( 'liveNotification' ); 
    this . logger . log ( `Heartbeat 消息发送到具有活动发射器侦听器计数的 SSE 客户端:${emitterListenerCount} ` ); 
    this . emitter . emit ( 'liveNotification' , { data : { content : 'heartbeat' } });
  }, 30000 ); 
}

性能表现

目前的通知系统有 90,000 个并发连接,在生产环境中运行 15 个 pod 来处理负载。每个 Kubernetes pod 消耗大约 800 MB 的内存和 300 Mi 的 CPU 资源。

相关推荐

Dubai's AI Boom Lures Global Tech as Emirate Reinvents Itself as Middle East's Silicon Gateway

AI-generatedimageAsianFin--Dubaiisrapidlytransformingitselffromadesertoilhubintoaglob...

OpenAI Releases o3-pro, Cuts o3 Prices by 80% as Deal with Google Cloud Reported to Make for Compute Needs

TMTPOST--OpenAIisescalatingthepricewarinlargelanguagemodel(LLM)whileseekingpartnershi...

黄仁勋说AI Agent才是未来!但究竟有些啥影响?

,抓住风口(iOS用户请用电脑端打开小程序)本期要点:详解2025年大热点你好,我是王煜全,这里是王煜全要闻评论。最近,有个词被各个科技大佬反复提及——AIAgent,智能体。黄仁勋在CES展的发布...

商城微服务项目组件搭建(五)——Kafka、Tomcat等安装部署

1、本文属于mini商城系列文档的第0章,由于篇幅原因,这篇文章拆成了6部分,本文属于第5部分2、mini商城项目详细文档及代码见CSDN:https://blog.csdn.net/Eclipse_...

Python+Appium环境搭建与自动化教程

以下是保姆级教程,手把手教你搭建Python+Appium环境并实现简单的APP自动化测试:一、环境搭建(Windows系统)1.安装Python访问Python官网下载最新版(建议...

零配置入门:用VSCode写Java代码的正确姿

一、环境准备:安装JDK,让电脑“听懂”Java目标:安装Java开发工具包(JDK),配置环境变量下载JDKJava程序需要JDK(JavaDevelopmentKit)才能运行和编译。以下是两...

Mycat的搭建以及配置与启动(mycat2)

1、首先开启服务器相关端口firewall-cmd--permanent--add-port=9066/tcpfirewall-cmd--permanent--add-port=80...

kubernetes 部署mysql应用(k8s mysql部署)

这边仅用于测试环境,一般生产环境mysql不建议使用容器部署。这里假设安装mysql版本为mysql8.0.33一、创建MySQL配置(ConfigMap)#mysql-config.yaml...

Spring Data Jpa 介绍和详细入门案例搭建

1.SpringDataJPA的概念在介绍SpringDataJPA的时候,我们首先认识下Hibernate。Hibernate是数据访问解决技术的绝对霸主,使用O/R映射(Object-Re...

量子点格棋上线!“天衍”邀您执子入局

你是否能在策略上战胜量子智能?这不仅是一场博弈更是一次量子智力的较量——量子点格棋正式上线!试试你能否赢下这场量子智局!游戏玩法详解一笔一画间的策略博弈游戏目标:封闭格子、争夺领地点格棋的基本目标是利...

美国将与阿联酋合作建立海外最大的人工智能数据中心

当地时间5月15日,美国白宫宣布与阿联酋合作建立人工智能数据中心园区,据称这是美国以外最大的人工智能园区。阿布扎比政府支持的阿联酋公司G42及多家美国公司将在阿布扎比合作建造容量为5GW的数据中心,占...

盘后股价大涨近8%!甲骨文的业绩及指引超预期?

近期,美股的AI概念股迎来了一波上升行情,微软(MSFT.US)频创新高,英伟达(NVDA.US)、台积电(TSM.US)、博通(AVGO.US)、甲骨文(ORCL.US)等多股亦出现显著上涨。而从基...

甲骨文预计新财年云基础设施营收将涨超70%,盘后一度涨8% | 财报见闻

甲骨文(Oracle)周三盘后公布财报显示,该公司第四财季业绩超预期,虽然云基建略微逊于预期,但管理层预计2026财年云基础设施营收预计将增长超过70%,同时资本支出继上年猛增三倍后,新财年将继续增至...

Springboot数据访问(整合MongoDB)

SpringBoot整合MongoDB基本概念MongoDB与我们之前熟知的关系型数据库(MySQL、Oracle)不同,MongoDB是一个文档数据库,它具有所需的可伸缩性和灵活性,以及所需的查询和...

Linux环境下,Jmeter压力测试的搭建及报错解决方法

概述  Jmeter最早是为了测试Tomcat的前身JServ的执行效率而诞生的。到目前为止,它的最新版本是5.3,其测试能力也不再仅仅只局限于对于Web服务器的测试,而是涵盖了数据库、JM...

取消回复欢迎 发表评论: