如何在数据库更新时发送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' }; }
行业标准实践
- 连接生命周期管理:必须监听客户端的
close事件,及时从连接池中移除无效连接,防止内存泄漏。 - 遵循SSE规范:消息格式必须以
data:开头,结尾用两个换行符\n\n分隔;推送JSON数据时,需用JSON.stringify转换为字符串。 - 错误处理:发送消息时捕获异常,处理客户端已断开但未触发
close事件的场景,及时清理无效连接。 - 认证与授权:若SSE端点涉及敏感数据,需在
/send-events中加入认证逻辑(如验证token、会话),仅允许合法客户端建立连接。 - 事件类型区分:若有多种更新场景,可在消息中添加
event:字段指定事件类型(如event: user-profile-update\n),客户端可监听特定事件,减少不必要的处理。 - 心跳机制:浏览器会自动断开长时间无数据的长连接,可定期发送心跳消息(如每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); - 限流与资源控制:限制同时连接的客户端数量,避免过多连接耗尽服务器资源;分布式场景下可结合Redis等工具实现跨节点的SSE推送。
内容的提问来源于stack exchange,提问作者Kallol Medhi
相关产品推荐
相关产品推荐

