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

摘要:在AI数据查询工具中,用户提问后SQL生成是流式的,但SQL执行可能耗时数秒甚至数十秒。如果让用户干等,体验直接崩掉;如果允许并发提问,又面临状态错乱、结果串台、任务丢失的风险。本文拆解我们如何用"独立会话 + 异步队列 + SSE回推"三件套解决这个问题:每个问题一个sessionId贯穿全链路,执行任务进Bull队列持久化,结果通过事件总线精准路由回对应会话的SSE连接。上线后日均处理3000+次查询,任务丢失率为零,用户可以边等结果边问下一个问题,全程无阻塞。

这是该系列的收官篇,前三篇分别讲了Class→Hooks重构、SSE流式引擎、可配置卡片系统,这篇补上最后一块拼图。


一、为什么不能"生成完就同步执行"?

最直觉的做法是:大模型吐完SQL → 后端立即查库 → 等结果回来再一次性返回给前端。这在demo里没问题,生产环境会炸:

  1. HTTP超时:复杂查询跑20s+很常见,Nginx默认60s超时,网关层还可能更短。同步等待极易触发504;
  2. 连接占用:每个查询占一个HTTP连接直到完成,并发100个慢查询就打满连接池;
  3. 用户体验差:用户盯着loading转圈干等,没法继续问下一个问题,对话流被强行打断;
  4. 容错脆弱:执行过程中网络抖动、服务重启,任务直接蒸发,用户只能重新提问。

所以必须拆成两段:生成是实时的(SSE流式),执行是异步的(后台队列)
但这引出了新的工程问题:异步之后,结果怎么回到正确的人手里?
多个问题并发时状态怎么隔离?服务挂了任务怎么不丢?


二、整体架构:三段式流水线

用户提问 ──▶ [生成阶段] ──▶ [提交阶段] ──▶ [执行阶段] ──▶ [回传阶段]
SSE流式 校验+入队 Bull Worker SSE/历史拉取
(token推送) (毫秒级响应) (Redis持久化) (精准路由)

每个阶段都有明确的边界和契约,通过 sessionId 串联。下面逐段拆解。


三、独立会话模型:每个问题一个"房间"

这是整个方案的基石。会话 ≠ 用户登录态,会话 = 一次完整的问答生命周期

3.1 sessionId的设计

// 格式:{userId}:{timestamp}:{randomSuffix}
// 示例:u_10086:1725984000000:a3f8k2
function createSessionId(userId: string): string {
return `${userId}:${Date.now()}:${nanoid(6)}`;
}

为什么不用UUID?因为sessionId需要携带业务语义:

  • userId 前缀方便按用户维度排查问题、做权限校验;
  • timestamp 让sessionId天然有序,日志检索和时间线还原不需要额外排序;
  • nanoid 后缀保证同一毫秒内的并发唯一性。

3.2 前端多会话状态管理

interface SessionState {
id: string;
status: 'streaming' | 'executing' | 'completed' | 'failed';
messages: Message[]; // 对话历史
sql?: string; // 生成的SQL
result?: QueryResult; // 查询结果
error?: string;
createdAt: number;
}

// 核心:Map结构,key=sessionId,各会话完全隔离
const sessions = useRef<Map<string, SessionState>>(new Map());

// 切换当前活跃会话 ≠ 销毁其他会话
function switchSession(sessionId: string) {
activeSessionId.current = sessionId;
// 其他会话的状态完整保留,SSE连接保持或按需重连
}

关键原则:切换会话只是UI层的焦点转移,后台的生成、执行、回传都不受影响。用户可以A问题还在跑的时候切到B问题继续聊,A的结果回来后自动更新到A会话里,不会弹到B界面上。


四、异步执行队列:Bull + Redis的工程细节

我们选Bull而不是自己写setTimeout/setInterval,原因很简单:持久化、重试、可观测性都是生产级刚需

4.1 任务入队

import { Queue } from 'bullmq';

const queryQueue = new Queue('sql-execution', { connection: redisConnection });

async function submitQuery(sessionId: string, sql: string, userId: string) {
await queryQueue.add(
'execute',
{ sessionId, sql, userId },
{
jobId: sessionId, // 用sessionId做jobId,天然幂等
attempts: 3, // 失败重试3次
backoff: { type: 'exponential', delay: 2000 },
removeOnComplete: 1000, // 成功后保留1000条记录用于排查
removeOnFail: 5000, // 失败保留更多
}
);
}

几个容易忽略的点:

  • jobId = sessionId:同一个会话重复提交不会创建新任务,避免重复执行。前端"重试"按钮调同样的submitQuery即可,无需额外去重逻辑。
  • removeOnComplete/Fail:不设的话Redis内存会无限增长。但也不能设太小,否则出问题时无迹可查。我们根据日均量估算保留窗口,兼顾内存和可追溯性。
  • 任务payload要小:只存sessionId和SQL,不存大对象。结果数据走单独的存储(见下文),队列只负责调度。

4.2 Worker执行与状态上报

import { Worker } from 'bullmq';

const worker = new Worker('sql-execution', async (job) => {
const { sessionId, sql, userId } = job.data;

// 1. 通知前端"开始执行"
await eventBus.publish(sessionId, {
event: 'exec_status',
data: { status: 'running', progress: 0 }
});

// 2. 执行SQL(带超时保护)
const result = await executeWithTimeout(sql, {
timeout: 30000,
maxRows: 10000,
userId, // 用于行级权限过滤
});

// 3. 结果落库(不通过SSE直传大数据)
const resultId = await storeResult(sessionId, result);

// 4. 通知前端"执行完成",只传resultId
await eventBus.publish(sessionId, {
event: 'result',
data: { resultId, rowCount: result.rows.length }
});

return { resultId }; // Bull记录返回值
}, { concurrency: 20 }); // 控制并发,防止打爆数据库

为什么结果不直接通过SSE推? 因为SSE是文本协议,大结果集序列化+传输可能耗时数秒,期间SSE连接被阻塞,心跳中断可能被客户端判定为断连。我们把结果存到Redis/Tair(TTL 1小时),SSE只推一个轻量的resultId,前端收到后再单独请求结果数据。生成通道和数据通道分离,互不干扰。

4.3 超时与资源保护

async function executeWithTimeout(sql: string, opts: TimeoutOpts) {
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), opts.timeout);

try {
return await db.query(sql, { signal: controller.signal, maxRows: opts.maxRows });
} catch (err) {
if (controller.signal.aborted) {
throw new Error(`Query timeout after ${opts.timeout}ms`);
}
throw err;
} finally {
clearTimeout(timer);
}
}

数据库侧也有兜底:查询账号绑定资源组,单查询CPU时间上限30s,超限自动kill。
应用层超时 + 数据库层限流双重保险,防止一个烂查询拖垮整个实例。


五、结果回传:SSE在线推,离线拉

这是"零丢失"的关键。用户可能在任务执行期间关掉页面、切换标签页、甚至退出登录。

5.1 在线场景:事件总线 → SSE

// 事件总线基于Redis Pub/Sub,跨进程/跨实例通信
class EventBus {
async publish(sessionId: string, message: SSEMessage) {
// 写入Redis Stream作为持久化备份
await redis.xadd(`session:${sessionId}:events`, '*', ...serialize(message));
// 同时Pub/Sub实时推送
await redis.publish(`channel:session:${sessionId}`, JSON.stringify(message));
}
}

// SSE Server订阅对应channel
redis.subscribe(`channel:session:${sessionId}`, (msg) => {
const sseRes = sessionManager.get(sessionId)?.res;
if (sseRes && !sseRes.writableEnded) {
sseRes.write(`event: ${msg.event}\ndata: ${JSON.stringify(msg.data)}\n\n`);
}
});

为什么Pub/Sub还要加Redis Stream? 因为Pub/Sub是fire-and-forget,订阅者不在线消息就丢了。
Stream提供消费者组语义,SSE重连后可以从上次消费位点续读,补齐断线期间的事件。

5.2 离线场景:历史事件拉取

// 用户重新打开页面时调用
app.get('/api/session/:sessionId/events', async (req, res) => {
const events = await redis.xrange(
`session:${req.params.sessionId}:events`,
req.query.lastEventId || '-', // 支持续传
'+',
'COUNT', 100 // 单次最多100条,分页拉取
);
res.json(events.map(deserialize));
});

前端初始化时检查本地缓存的lastEventId,如果有未消费的事件则拉取补齐。
在线推 + 离线拉组合,覆盖了所有连接状态,这就是"零丢失"的工程保障。

5.3 结果数据的懒加载

// 前端收到result事件后
async function onResult(sessionId: string, resultId: string) {
// 不立即拉取,等用户切换到该会话或展开结果区域时才请求
const result = await fetch(`/api/result/${resultId}`).then(r => r.json());
renderChartOrTable(sessionId, result);
}

日均3000次查询,但同一时刻用户真正在看的结果可能只有几十个。
懒加载避免了无效的数据传输和渲染开销。


六、状态机:让会话生命周期可预测

每个会话有明确的状态流转,杜绝"不知道当前在干嘛"的混沌状态:

┌──────────┐
│ idle │ ◀──────────────────────┐
└────┬─────┘ │
│ 用户提问 │ 完成/失败/超时
▼ │
┌──────────┐ │
│streaming │ ──── SQL生成完成 ──────▶ │
└────┬─────┘ │
│ 校验通过 │
▼ │
┌──────────┐ │
│ queued │ ──── Worker拾取 ──────▶ │
└────┬─────┘ │
│ │
▼ │
┌──────────┐ │
│ running │ ──── 执行完成/失败 ────▶ │
└──────────┘

状态变更全部通过事件驱动,前端监听事件更新UI,不做任何隐式状态推断。调试时只要看事件序列就能还原完整生命周期,比翻代码快10倍。

我们还加了状态一致性校验:Worker执行前检查会话是否仍处于queued状态(防止重复提交导致重复执行),结果回传前检查会话是否仍存在(防止已清理的会话收到脏事件)。每一跳都有守卫,状态机才不会跑飞


七、可观测性:出了问题能快速定位

日均3000次查询,没有监控就是盲人摸象。我们埋了三层:

  1. 链路追踪:sessionId作为traceId贯穿BFF → LLM → Queue → Worker → DB,Jaeger里一条链路看全貌;
  2. 队列指标:Grafana面板实时展示队列长度、平均等待时间、执行P99延迟、失败率。设置告警:队列积压>50或P99>15s自动钉钉通知;
  3. 会话级日志:每个sessionId对应一个结构化日志条目,记录状态变更、事件收发、耗时分布。运营反馈"我的查询没出结果"时,拿sessionId搜日志,30秒内定位原因。

可观测性不是锦上添花,是异步系统的生命线。同步调用栈断了看报错就行,异步链路断了没有日志就是黑洞。


八、踩坑备忘录

  1. Bull的concurrency别设太高:我们最初设50,结果数据库连接池被打满,反而整体变慢。压测后定为20,配合连接池max=30,吞吐和稳定性达到平衡。
  2. Redis Stream的MAXLEN要设:不设的话Stream无限增长。我们设MAXLEN ~ 200,足够覆盖单次会话的全部事件,又不会撑爆内存。
  3. SSE重连后的事件去重:Stream续传可能重发已消费的事件。前端维护一个processedEventIds Set,重复事件静默跳过。
  4. Worker进程崩溃时的任务恢复:Bull默认会把active状态的job标记为stalled并重新入队。但要确保SQL执行是幂等的(SELECT天然幂等,如果是写操作就需要额外设计)。
  5. 前端内存泄漏:多会话场景下,已完成的会话如果不主动清理,Map会无限增长。我们做了LRU策略:保留最近20个活跃会话,超出的自动归档到IndexedDB,需要时再恢复。
  6. 时区问题:sessionId里的timestamp是UTC,但运营看的是北京时间。所有展示层统一做时区转换,存储层永远用UTC。这个看似基础,但在跨时区协作和日志对齐时能省无数麻烦。

九、系列回顾与总结

四篇文章串起来,是一套完整的AI数据查询工具技术栈:

篇目 解决的问题 核心技术点
Class→Hooks重构 遗留代码维护成本高 AI辅助迁移Prompt策略、人机协作校验
SSE流式引擎 大模型输出体验卡顿 Node.js SSE实现、增量渲染优化、SQL校验
可配置卡片系统 查询条件扩展成本高 JSON Schema驱动、渲染引擎抽象、AI回填映射
异步执行与会话管理 并发查询状态混乱、任务丢失 独立会话模型、Bull队列、事件总线、在线推+离线拉

回头看,这套方案没有用什么前沿黑科技。SSE、Bull、Redis Stream、React Hooks都是成熟技术,真正的功夫花在"把它们正确地组装在一起"上:接口契约清晰、状态边界明确、容错层层兜底、可观测性贯穿始终。

这也是我一直想传达的观点:工程能力的体现不在于用了多新的技术,而在于面对真实约束时,能否用合适的技术组合出稳定、可维护、可复用的解决方案
这套模式后来被团队3个项目复用,不是因为代码写得"通用",而是因为抽象层次恰好卡在业务概念的正确位置上。

希望这个系列对你有参考价值。如果有任何问题或想讨论的细节,欢迎评论区交流。


后续会持续分享AI工程化、前端效能提升相关的实战内容,收藏网址不迷路。

评论