如何在Express中将MSSQL查询结果流式传输至HTTP响应
嘿,我来帮你搞定把mssql查询结果流式传输到Express HTTP响应的问题!你已经搞定了命令行的流式输出,只需要把输出目标从控制台转到Express的res对象就好,下面是具体的实现方案:
完整实现代码
首先,我们先修改你的index.js中的路由逻辑,把流式数据直接写入响应:
const sql = require('mssql'); const config = require('./db').config; // 假设你的config在db.js导出 const router = require('express').Router(); router.get('/', async function (req, res, next) { let pool; try { // 初始化连接池(后面会说更高效的复用方式) pool = await new sql.ConnectionPool(config).connect(); const request = new sql.Request(pool); // 开启mssql的流式模式 request.stream = true; // 设置响应头,这里我们输出标准JSON数组格式 res.setHeader('Content-Type', 'application/json'); // 先写入数组的开头括号 res.write('['); let isFirstRow = true; // 监听每一行数据的事件 request.on('row', (row) => { // 非第一行需要加逗号分隔JSON对象 if (!isFirstRow) { res.write(','); } isFirstRow = false; // 将行数据转为JSON字符串写入响应 const writeSuccess = res.write(JSON.stringify(row)); // 处理背压:如果响应缓冲区满了,暂停读取数据库流 if (!writeSuccess) { request.pause(); res.once('drain', () => { request.resume(); }); } }); // 监听查询完成事件,关闭JSON数组并结束响应 request.on('requestCompleted', () => { res.write(']'); res.end(); // 关闭本次请求创建的连接池 pool.close(); }); // 监听查询错误,及时结束响应并释放资源 request.on('error', (err) => { console.error('流式查询出错:', err); // 确保响应被正确结束,避免客户端挂起 if (!res.headersSent) { res.status(500).json({ error: '数据查询失败' }); } else { res.end(); } pool.close(); }); // 执行你的查询语句,替换成实际的SQL await request.query('SELECT * FROM YourLargeDatasetTable'); } catch (err) { console.error('数据库连接失败:', err); if (!res.headersSent) { res.status(500).json({ error: '数据库连接异常' }); } if (pool) pool.close(); } }); module.exports = router;
关键优化点说明
1. 连接池复用(重要!)
上面的代码每次请求都新建连接池,效率很低,建议在db.js中初始化全局连接池复用:
修改db.js:
const sql = require('mssql'); const config = { user: 'user', password: 'pass', server: 'host', database: 'db' }; // 创建全局连接池Promise,只初始化一次 const poolPromise = new sql.ConnectionPool(config) .connect() .then(pool => { console.log('数据库连接成功'); return pool; }) .catch(err => console.error('数据库连接失败:', err)); module.exports = { sql, poolPromise };
然后在index.js中使用:
const { poolPromise } = require('./db'); router.get('/', async function (req, res, next) { try { // 直接复用全局连接池 const pool = await poolPromise; const request = new sql.Request(pool); // 后面的流式逻辑和之前一致,不需要再新建连接池 // ... } catch (err) { // 错误处理逻辑 } });
2. 可选的响应格式:JSON Lines
如果你的数据量特别大,不想构建完整的JSON数组,可以用JSON Lines格式(每行一个JSON对象),这种格式客户端可以逐行解析,内存压力更小:
修改响应头和写入逻辑:
// 设置JSON Lines的响应头 res.setHeader('Content-Type', 'application/x-ndjson'); // row事件中直接写入每行JSON+换行 request.on('row', (row) => { const writeSuccess = res.write(JSON.stringify(row) + '\n'); // 同样处理背压 if (!writeSuccess) { request.pause(); res.once('drain', () => request.resume()); } }); // 完成时直接结束响应,不需要括号 request.on('requestCompleted', () => { res.end(); });
3. 背压处理
代码中加入了背压处理逻辑:当res.write()返回false时,说明响应缓冲区已满,此时暂停读取数据库的流,等res的drain事件触发(缓冲区清空)后再恢复读取,避免内存溢出。
内容的提问来源于stack exchange,提问作者Blair
相关产品推荐
相关产品推荐

