抱歉,您的浏览器无法访问本站
本页面需要浏览器支持(启用)JavaScript
了解详情 >

摘要:运营同学想看个数据,提需求、排期、开发、上线,一套流程走下来至少一周。我们基于 Node.js + SSE 搭建了一个流式SQL生成引擎,让运营通过对话直接取数:输入自然语言问题,大模型流式输出SQL构建过程,提交后异步执行,结果自动渲染成图表或表格。本文还原整个工具从0到1的技术落地过程,重点讲SSE的前后端实现细节、长文本增量渲染优化,以及如何对AI生成的SQL做基础校验,避免"垃圾进、垃圾出"。

上线后日均处理3000+次查询,几乎零任务丢失。这套模式后来被团队3个AI项目直接复用。


一、业务背景:为什么选SSE而不是WebSocket?

先说场景。运营的日常诉求很简单:"帮我查一下上周华东区各品类的销售额"。
传统流程是提需求 → 数据开发排期 → 写SQL → 出结果,快则三天慢则一周。
我们想做的工具是:运营自己对话,实时看到SQL是怎么一步步拼出来的,确认后提交执行,结果自动出图表。

技术选型上有三个候选:轮询、WebSocket、SSE。最终选SSE,理由很实际:

  1. 单向场景:大模型生成是纯服务端→客户端的推送,不需要双向通信,WebSocket的长连接维护、心跳、重连机制都是多余的成本;
  2. 协议亲和:SSE基于HTTP,天然穿透公司网关、Nginx代理和鉴权中间件,运维侧几乎零改造;
  3. 自带容错:浏览器原生 EventSource 断线自动重连,配合Last-Event-ID可以续传。

一句话:大模型流式输出这个场景,SSE是性价比最高的方案。这也是为什么OpenAI等厂商的API都默认用SSE。


二、整体架构

用户提问


┌─────────────┐ ┌──────────────┐ ┌─────────────┐
│ Node.js BFF │───▶│ LLM API │ │ SQL校验层 │
│ (SSE Server)│◀───│ (流式返回) │ │ 语法+关键字 │
└──────┬──────┘ └──────────────┘ └──────┬──────┘
│ SSE: token流 / 执行状态 │
▼ ▼
┌─────────────┐ ┌─────────────┐
│ 前端增量渲染 │ │ 异步执行队列 │
└─────────────┘ │ (每个问题独立会话)│
└─────────────┘

核心流程分两段:

  • 第一段(生成):用户提问 → BFF转发给大模型 → 模型流式返回SQL构建过程 → SSE推给前端逐字渲染;
  • 第二段(执行):SQL通过校验后提交异步任务 → 前端保持SSE连接接收执行状态 → 结果回来后渲染图表/表格。用户此时可以继续问下一个问题,每个问题一个独立会话,互不阻塞

三、服务端实现:Node.js + SSE的三个关键点

3.1 基础的SSE响应头与事件推送

// Node.js + Express 示例
app.get('/api/chat/stream', async (req, res) => {
// 关键响应头:三个一个都不能少
res.writeHead(200, {
'Content-Type': 'text/event-stream; charset=utf-8',
'Cache-Control': 'no-cache, no-transform',
'Connection': 'keep-alive',
'X-Accel-Buffering': 'no', // 关键!禁用Nginx缓冲,否则数据会被攒批下发
});

const sessionId = req.query.sessionId;
sessionManager.register(sessionId, res);

try {
// 对接大模型的流式API
const llmStream = await callLLMStream({
prompt: buildPrompt(req.query.question),
});

for await (const chunk of llmStream) {
// SSE标准格式:data字段 + 两个换行
res.write(`event: message\n`);
res.write(`id: ${++eventSeq}\n`);
res.write(`data: ${JSON.stringify({ type: 'token', content: chunk.text })}\n\n`);
}

res.write(`event: done\ndata: ${JSON.stringify({ type: 'complete' })}\n\n`);
} catch (err) {
res.write(`event: error\ndata: ${JSON.stringify({ message: err.message })}\n\n`);
} finally {
sessionManager.unregister(sessionId);
res.end();
}

// 客户端断开时清理资源
req.on('close', () => sessionManager.unregister(sessionId));
});

三个容易踩的坑:

  1. **X-Accel-Buffering: no**:不加这个头,Nginx会缓冲响应体,前端表现为"卡半天然后一次性吐出来",流式效果完全失效。这是SSE过代理最常见的坑,没有之一。
  2. 心跳保活:长连接空闲超过60s可能被网关掐断。我们每15s发一条注释帧 : ping\n\n(冒号开头的行是SSE注释,浏览器会忽略但能保活连接)。
  3. req.on('close') 清理:用户关掉页面时如果不注销会话,res对象会泄漏,高并发下内存持续上涨。

3.2 事件分类设计:一个通道传多种消息

生成过程和执行状态都走同一条SSE连接,靠 event 字段区分类型:

event类型 含义 前端处理
message 模型输出的token流 追加渲染到对话气泡
sql_ready 完整SQL已生成 展示SQL预览+确认按钮
exec_status 异步任务状态(排队/执行中/完成) 更新进度指示器
result 查询结果数据 渲染图表或表格
error 错误信息 提示并允许重试
done 本轮流结束 关闭loading态

这个设计让一条连接承载了完整的会话生命周期,前端不需要为"生成"和"执行"维护两套连接逻辑。

3.3 会话管理:每个问题独立,支持并发提问

这是"运营可以边等结果边问下一个问题"的实现基础:

class SessionManager {
constructor() {
// sessionId -> { res, abortController, status }
this.sessions = new Map();
}

register(sessionId, res) {
this.sessions.set(sessionId, {
res,
abortController: new AbortController(),
status: 'streaming',
createdAt: Date.now(),
});
}

// 用户切换问题 ≠ 中断旧会话
// 旧会话转入后台继续执行,完成后通过事件总线推送结果
detach(sessionId) {
const s = this.sessions.get(sessionId);
if (s) s.status = 'background';
}
}

每个问题生成唯一的 sessionId,前端用 Map<sessionId, 消息列表> 维护多会话状态。
SQL提交后任务进入异步队列(我们用的Bull + Redis),执行结果通过事件总线路由回对应会话的SSE连接。
上线后日均3000+次查询,任务丢失率为零——关键在于每个环节(生成、提交、执行、回传)都有sessionId贯穿 + 状态持久化,任何一步断掉都能靠 exec_status 事件对账恢复。


四、前端实现:增量渲染,长回复不卡顿

4.1 用fetch + ReadableStream代替EventSource

浏览器原生 EventSource 有个硬伤:只支持GET,不能带自定义Header(比如鉴权token)。我们的做法是用 fetch 读流:

async function streamChat(sessionId: string, question: string) {
const response = await fetch(`/api/chat/stream?sessionId=${sessionId}&question=${encodeURIComponent(question)}`, {
headers: { Authorization: `Bearer ${token}` },
signal: abortController.signal, // 支持手动中断
});

const reader = response.body!.getReader();
const decoder = new TextDecoder();
let buffer = '';

while (true) {
const { done, value } = await reader.read();
if (done) break;

buffer += decoder.decode(value, { stream: true });
// 按SSE协议边界切分:一个完整事件以\n\n结尾
const events = buffer.split('\n\n');
buffer = events.pop()!; // 最后一段可能不完整,留在buffer等下一chunk

for (const raw of events) {
const event = parseSSE(raw); // 解析event/data字段
dispatchEvent(sessionId, event);
}
}
}

注意 buffer 的处理:网络分包不会按 \n\n 边界切,半个JSON直接parse会崩,必须缓存拼接。

4.2 token级渲染的性能陷阱与解法

大模型每秒可能吐50+个token,如果每个token都触发一次React setState + 重渲染,长回复(几千字符的SQL+解释)会让页面明显掉帧。我们的三层优化:

① 批量合并(Buffer + rAF)

// token先攒进ref,不直接setState
const tokenBuffer = useRef('');
const rafId = useRef(0);

function onToken(text: string) {
tokenBuffer.current += text;
if (!rafId.current) {
rafId.current = requestAnimationFrame(() => {
setContent(prev => prev + tokenBuffer.current);
tokenBuffer.current = '';
rafId.current = 0;
});
}
}

把"N次setState"合并成"每帧一次",渲染频率从50+/s降到60帧对齐的≤60/s,且和浏览器绘制节奏一致。

② 流式内容组件隔离

正在流式输出的气泡拆成独立组件 StreamingBubble,用 React.memo 包裹,token更新只重渲染这一个气泡,历史消息列表完全不动。

③ 超长内容分段

回复超过2000字符后,把已完成的部分"冻结"成静态节点,只有尾部活跃段参与流式更新。这样即使模型输出上万字符,每帧的diff成本也恒定。

实测:优化前渲染5000字符回复时帧率掉到20fps以下,优化后全程稳定55fps+。

4.3 SQL高亮的懒加载

生成过程中SQL代码块需要语法高亮,但highlight.js对每个token都重新高亮整个代码块的开销很大。
解法:流式期间纯文本渲染,sql_ready 事件到达后(即内容完整了)才做一次高亮
用户在逐字输出阶段本来也看不清代码细节,体验无损,性能大增。


五、SQL校验层:AI的输出不能裸奔进数据库

这是整个系统的安全底线。大模型生成的SQL可能语法错误、可能全表扫描、极端情况下甚至可能带危险语句。
我们做了双层拦截,在Node.js侧执行(不放前端,防止绕过):

5.1 关键字黑名单过滤

const FORBIDDEN_PATTERNS = [
/\b(DROP|TRUNCATE|ALTER|CREATE)\b/i, // DDL全禁
/\b(DELETE|UPDATE|INSERT|REPLACE)\b/i, // DML全禁,只允许查询
/\bGRANT\b|\bREVOKE\b/i, // 权限操作
/;\s*\S/, // 多语句注入(分号后还有内容)
/\bINTO\s+OUTFILE\b|\bLOAD_FILE\b/i, // 文件读写
/\/\*[\s\S]*?\*\//, // 注释内藏语句
];

function keywordCheck(sql) {
for (const pattern of FORBIDDEN_PATTERNS) {
if (pattern.test(sql)) {
return { pass: false, reason: `命中危险模式: ${pattern}` };
}
}
return { pass: true };
}

5.2 语法解析 + 结构校验

光靠正则不够(比如 DELETE 出现在字符串字面量里会误杀),所以我们引入 node-sql-parser 做AST级校验:

const { Parser } = require('node-sql-parser');
const parser = new Parser();

function syntaxCheck(sql) {
try {
const ast = parser.astify(sql);
const statements = Array.isArray(ast) ? ast : [ast];

// 规则1:必须是单条SELECT语句
if (statements.length !== 1 || statements[0].type !== 'select') {
return { pass: false, reason: '仅允许单条SELECT查询' };
}
// 规则2:必须带LIMIT,防止拖库
if (!statements[0].limit) {
return { pass: false, reason: '查询必须包含LIMIT', suggestion: '自动追加 LIMIT 10000' };
}
// 规则3:表白名单校验(解析出的表必须在授权的元数据范围内)
const tables = parser.tableList(sql).map(t => t.split('`')[1]);
const illegal = tables.filter(t => !TABLE_WHITELIST.has(t));
if (illegal.length) {
return { pass: false, reason: `无权访问的表: ${illegal.join(', ')}` };
}
return { pass: true };
} catch (e) {
return { pass: false, reason: `语法错误: ${e.message}` };
}
}

5.3 校验失败怎么办?回喂给模型自修复

校验不通过时不是直接报错给用户,而是把错误原因拼回Prompt让模型重新生成(最多重试2次):

你上次生成的SQL未通过校验,原因:查询必须包含LIMIT。
请修正后重新输出完整SQL。用户原始问题:{question}
上次生成的SQL:{sql}

实践中约85%的校验失败能通过一轮自修复解决,用户侧几乎无感知。
仍失败则降级为提示"该问题暂时无法自动取数,已通知数据团队",保证体验兜底。

此外,执行层还有最后一道保险:查询账号本身就是只读权限 + 资源组限流(单查询超时30s自动kill),即使前三层全部失效,也造不成实际破坏。纵深防御,每层都不指望上一层是完美的。


六、踩坑备忘录

  1. Nginx缓冲X-Accel-Buffering: no + proxy_buffering off,双重确认。这个坑排查起来极其隐蔽——本地好好的,一上测试环境流式就变"批式"。
  2. TextDecoder的stream参数decoder.decode(value, { stream: true }) 必须传,否则多字节中文字符被chunk从中间截断时会出乱码。
  3. EventSource重连风暴:如果用原生EventSource,服务端异常关闭会触发浏览器疯狂重连。改用fetch方案后自己控制重试策略(指数退避,最多3次)。
  4. 异步结果回推时连接已断:用户关了页面,任务还在跑。解法是结果落库 + 用户下次打开时拉取历史会话,"任务丢失"感知归零。
  5. Prompt里放表结构要精简:早期把全量表结构塞进Prompt,token爆炸且模型注意力涣散。后来改成根据问题关键词检索Top-5相关表的schema注入,SQL准确率从72%提到91%。

七、写在最后

回头看,这个项目技术上没有特别"炫"的部分——SSE是成熟协议,SQL校验是常规手段,增量渲染是标准优化。
它真正的价值在于把一串朴素的技术正确地组装起来,解决了真实的人效问题
运营从"提需求等排期"变成"自己对话取数",单次取数耗时从平均3天降到分钟级。

也正因为方案简单、边界清晰,这套"SSE对接 + 卡片交互 + 异步执行"的模式后来被团队3个AI项目直接复用。
好的架构不是复杂的架构,是能被别人抄作业的架构。


下一篇打算写那个被复用最多的部分——可配置卡片选择系统(20+种查询卡片,JSON配置驱动,新业务接入零代码)。感兴趣的话欢迎收藏网址。

评论