颜力刚Ligang Yan

用 Server-Sent Events 从 NestJS 向 Next.js 实时推送

一步一步的教程:在 NestJS 里用 RxJS Subject 暴露一个 SSE 端点,在 Next.js 里用一个小小的 useSse hook 消费它,然后从后端任何服务推送事件。数据只往一个方向流的时候,比 WebSocket 简单得多。

nestjsnextjsssetutorial

English version: Real-time updates from NestJS to Next.js with Server-Sent Events

这篇教程带你实现一条实时通道,让 NestJS 后端主动向 Next.js 前端推送数据。用的是 Server-Sent Events(SSE),一种简单有效的单向服务器到客户端通信技术,直接建在 HTTP 之上。

适合实时通知、状态更新、给仪表盘喂实时数据这类场景。

第一部分:后端(NestJS)

目标是在 api 应用里建一个专门的端点,客户端订阅它就能收到一条事件流。

第 1.1 步:创建 Events 模块

先建一个 events 模块,代码归置整齐。

mkdir -p apps/api/src/modules/events

第 1.2 步:创建 Events 服务

这个服务用 RxJS 的 Subject 管理事件流,它是所有消息的中心通道。应用里其他部分通过这个服务发送新事件。

创建 apps/api/src/modules/events/events.service.ts:

import { Injectable } from '@nestjs/common';
import { Observable, Subject } from 'rxjs';

export interface MessageEvent {
  data: string | object;
}

@Injectable()
export class EventsService {
  private readonly events = new Subject<MessageEvent>();

  // 控制器调用它获取事件流
  getEvents(): Observable<MessageEvent> {
    return this.events.asObservable();
  }

  // 其他服务调用它向流里推一条新事件
  emitEvent(event: MessageEvent) {
    this.events.next(event);
  }
}

第 1.3 步:创建 Events 控制器

控制器暴露对外的端点。需要一个供客户端连接的(/sse),再加一个方便测试的辅助端点(/emit)。

创建 apps/api/src/modules/events/events.controller.ts:

import { Controller, Sse, Post, Body } from '@nestjs/common';
import { EventsService, MessageEvent } from './events.service';
import { Observable } from 'rxjs';
import { map } from 'rxjs/operators';

@Controller('events')
export class EventsController {
  constructor(private readonly eventsService: EventsService) {}

  @Sse('sse')
  sse(): Observable<MessageEvent> {
    return this.eventsService.getEvents().pipe(
      // SSE 的数据必须是 { data: your_data } 这个格式
      map((event: MessageEvent) => ({ data: event.data }))
    );
  }

  @Post('emit')
  emitEvent(@Body() body: { message: string }) {
    this.eventsService.emitEvent({ data: { content: body.message, timestamp: new Date() } });
    return { success: true };
  }
}

第 1.4 步:定义 Events 模块

把服务和控制器绑到一个模块里。

创建 apps/api/src/modules/events/events.module.ts:

import { Module } from '@nestjs/common';
import { EventsController } from './events.controller';
import { EventsService } from './events.service';

@Module({
  controllers: [EventsController],
  providers: [EventsService],
  // 导出服务,让应用里其他模块可以注入它并发送事件
  exports: [EventsService],
})
export class EventsModule {}

第 1.5 步:在应用里启用模块

把 EventsModule 导入主 AppModule。

修改 apps/api/src/app.module.ts:

import { Module } from '@nestjs/common';
import { ConfigModule } from '@nestjs/config';
import { AppController } from './app.controller';
import { AppService } from './app.service';
// 导入新模块
import { EventsModule } from './modules/events/events.module';

@Module({
  imports: [
    ConfigModule.forRoot({ isGlobal: true }),
    // ...你的其他模块
    // 把新模块加进 imports
    EventsModule,
  ],
  controllers: [AppController],
  providers: [AppService],
})
export class AppModule {}

后端完成。 API 现在可以推送事件了。


第二部分:前端(Next.js)

接下来配置 web 应用监听这些事件并显示出来。

第 2.1 步:写一个可复用的 SSE hook

一个自定义 React hook 最适合管理 EventSource 连接的生命周期。

先建 hooks 目录:

mkdir -p apps/web/hooks

创建 apps/web/hooks/use-sse.ts:

import { useState, useEffect } from 'react';

// 一个通用的 Server-Sent Events hook
export function useSse<T>(url: string, initialValue: T) {
  const [data, setData] = useState<T>(initialValue);

  useEffect(() => {
    // URL 指向 NestJS 的 SSE 端点
    const eventSource = new EventSource(url);

    eventSource.onopen = () => {
      console.log('SSE connection opened.');
    };

    eventSource.onmessage = (event) => {
      try {
        const parsedData = JSON.parse(event.data);
        setData(parsedData);
      } catch (error) {
        console.error('Failed to parse SSE data:', error);
      }
    };

    eventSource.onerror = (error) => {
      console.error('SSE error:', error);
      // 大多数错误浏览器会自动重连。
      // 想在特定错误上停止重试,可以在这里关闭。
      eventSource.close();
    };

    // 组件卸载时清理
    return () => {
      console.log('Closing SSE connection.');
      eventSource.close();
    };
  }, [url]); // URL 变化时重新执行

  return data;
}

第 2.2 步:创建 UI 组件

写一个组件,用这个 hook 显示实时数据,并带一个测试按钮。

创建 apps/web/components/real-time-updates.tsx:

'use client';

import { useSse } from '@/hooks/use-sse';
import { useState } from 'react';

// 定义期望从服务端收到的数据结构
interface SseData {
  content: string;
  timestamp: string;
}

export function RealTimeUpdates() {
  // 重要:换成你实际的 API 地址,最好放在环境变量里(比如 .env.local)
  const API_URL = process.env.NEXT_PUBLIC_API_URL || 'http://localhost:3001';
  const sseData = useSse<SseData | null>(`${API_URL}/api/events/sse`, null);
  const [status, setStatus] = useState('');

  const sendMessage = async () => {
    setStatus('Sending message...');
    try {
      const response = await fetch(`${API_URL}/api/events/emit`, {
        method: 'POST',
        headers: { 'Content-Type': 'application/json' },
        body: JSON.stringify({ message: `Hello from client at ${new Date().toLocaleTimeString()}` }),
      });
      if (response.ok) {
        setStatus('Message sent successfully! Waiting for SSE push...');
      } else {
        const error = await response.text();
        setStatus(`Failed to send message: ${error}`);
      }
    } catch (error) {
      setStatus(`Error: ${(error as Error).message}`);
    }
  };

  return (
    <div className="mt-8 p-4 border rounded-lg bg-gray-50 w-full">
      <h2 className="text-xl font-semibold mb-3">Real-Time Updates (via SSE)</h2>
      <div className="mb-4">
        <button
          onClick={sendMessage}
          className="px-4 py-2 bg-blue-500 text-white rounded hover:bg-blue-600 transition-colors"
        >
          Trigger Server Push
        </button>
        {status && <p className="text-sm text-gray-600 mt-2">{status}</p>}
      </div>
      <div className="p-3 bg-white border rounded shadow-inner min-h-[50px]">
        <p className="font-mono text-sm">
          <strong>Received Message:</strong>
          {sseData ?
            <span className="ml-2">{sseData.content} (@ {new Date(sseData.timestamp).toLocaleTimeString()})</span> :
            ' Waiting for server event...'
          }
        </p>
      </div>
    </div>
  );
}

第 2.3 步:把组件放进页面

把 RealTimeUpdates 放到任何你想要的页面。比如首页:

import { RealTimeUpdates } from '@/components/real-time-updates';

export default function HomePage() {
  return (
    <main className="flex flex-col items-center gap-8 px-4 py-20">
      <h1 className="text-3xl font-semibold">Home</h1>
      {/* 实时面板 */}
      <RealTimeUpdates />
    </main>
  );
}

第三部分:运行和测试

  1. 启动两个应用:分别用各自的 dev 脚本启动 NestJS API 和 Next.js 前端。
  2. 检查 API 地址:确认 real-time-updates.tsx 里的 API_URL(默认 http://localhost:3001)指向正在运行的 NestJS API。
  3. 打开首页:在浏览器里打开 Next.js 应用,应该能看到 “Real-Time Updates” 面板。
  4. 测试推送:点击 “Trigger Server Push” 按钮。状态文字会变化,随后 “Received Message” 会更新为服务端推过来的数据。

第四部分:接入真实业务

测试按钮很方便,但真实应用里事件应该由业务逻辑触发。因为 EventsService 已经导出,可以注入到任何其他服务里。

例子:下载完成时推一条通知。

// 在某个其他服务里,比如 downloader.service.ts
import { Injectable } from '@nestjs/common';
import { EventsService } from '../events/events.service';

@Injectable()
export class DownloaderService {
  // 在构造函数里注入 EventsService
  constructor(private readonly eventsService: EventsService) {}

  async processDownload(url: string) {
    // ... 下载文件的逻辑 ...
    console.log('Download complete.');

    // 向客户端推一条通知
    const notification = {
      message: `Finished downloading from ${url}`,
      timestamp: new Date(),
    };
    this.eventsService.emitEvent({ data: notification });
  }
}

到这里,你已经有了一个稳固、可扩展的基础,可以往应用里加各种实时功能。

这篇是我做一个下载工具时写的,后端需要在长任务完成时通知浏览器。SSE 正好够用:单向、纯 HTTP、浏览器自动重连、不用跑 WebSocket 服务器。

English version.