You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在数据库更新时发送Server-Sent Events(SSE)?

实现数据库更新触发SSE事件的方案及行业实践

核心思路

SSE(Server-Sent Events)是基于HTTP长连接的单向推送机制,要实现仅在数据库更新时推送事件,核心是维护活跃客户端连接池:当客户端连接SSE端点时,将其加入连接池;当数据库更新完成后,遍历连接池向所有活跃客户端发送最新数据。

代码实现步骤

1. 维护客户端连接池

定义一个全局集合存储所有活跃的SSE客户端响应对象,用于后续批量推送:

// 用Set存储活跃连接,自动去重且便于删除操作
const sseClients = new Set();

2. 修改SSE端点代码

调整/send-events端点,将客户端加入连接池,并处理连接断开的清理逻辑:

app.get('/send-events', (req, res) => {
    const headers = {
        Connection: "keep-alive",
        "Content-Type": "text/event-stream",
        "Cache-Control": "no-cache",
        // 跨域场景需配置,根据实际业务调整允许的源
        "Access-Control-Allow-Origin": "*"
    };
    res.writeHead(200, headers);

    // 发送初始连接确认(可选,告知客户端连接成功)
    res.write('data: connected\n\n');

    // 将当前客户端连接加入集合
    sseClients.add(res);

    // 监听客户端断开事件,移除无效连接避免内存泄漏
    req.on('close', () => {
        sseClients.delete(res);
        console.log('SSE client disconnected');
    });
});

3. 修改数据库更新API

在/update-db中完成数据库更新后,向所有活跃客户端推送最新数据:

app.post('/update-db', async (req, res) => {
    try {
        // 执行实际数据库更新操作,替换为你的业务逻辑
        const latestData = await updateDatabase(req.body);

        // 构造符合SSE规范的消息:data字段开头,结尾两个换行符
        const sseMessage = `data: ${JSON.stringify(latestData)}\n\n`;

        // 遍历所有活跃客户端发送消息
        sseClients.forEach(client => {
            try {
                client.write(sseMessage);
            } catch (error) {
                // 发送失败说明连接已失效,从集合中移除
                sseClients.delete(client);
                console.error('Failed to send SSE message:', error);
            }
        });

        res.status(200).send('db-updated');
    } catch (error) {
        console.error('Database update failed:', error);
        res.status(500).send('db-update-failed');
    }
});

// 模拟数据库更新函数,替换为你的实际数据库操作
async function updateDatabase(updateData) {
    // 示例:执行更新后返回最新数据
    return {
        updateTime: new Date().toISOString(),
        content: updateData,
        status: 'updated'
    };
}

行业标准实践

  1. 连接生命周期管理:必须监听客户端的close事件,及时从连接池中移除无效连接,防止内存泄漏。
  2. 遵循SSE规范:消息格式必须以data: 开头,结尾用两个换行符\n\n分隔;推送JSON数据时,需用JSON.stringify转换为字符串。
  3. 错误处理:发送消息时捕获异常,处理客户端已断开但未触发close事件的场景,及时清理无效连接。
  4. 认证与授权:若SSE端点涉及敏感数据,需在/send-events中加入认证逻辑(如验证token、会话),仅允许合法客户端建立连接。
  5. 事件类型区分:若有多种更新场景,可在消息中添加event: 字段指定事件类型(如event: user-profile-update\n),客户端可监听特定事件,减少不必要的处理。
  6. 心跳机制:浏览器会自动断开长时间无数据的长连接,可定期发送心跳消息(如每30秒发送data: heartbeat\n\n)维持连接活跃:
    // 30秒发送一次心跳
    setInterval(() => {
        const heartbeat = 'data: heartbeat\n\n';
        sseClients.forEach(client => {
            try {
                client.write(heartbeat);
            } catch (err) {
                sseClients.delete(client);
            }
        });
    }, 30000);
    
  7. 限流与资源控制:限制同时连接的客户端数量,避免过多连接耗尽服务器资源;分布式场景下可结合Redis等工具实现跨节点的SSE推送。

内容的提问来源于stack exchange,提问作者Kallol Medhi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 13:30:41