如何在服务端用pg-query-stream获取PostgreSQL数据?报错求助
问题描述
尝试在服务端通过pg-query-stream从PostgreSQL获取数据时,抛出错误:
Error: Invalid response from route /api: handler should return a Response object
推测原因是返回的是对象而非Response实例。当前代码如下:
服务端代码(routes/api/+server.ts)
/** @type {import('./$types').RequestHandler} */ export async function GET({ url }) { ... var server = http.createServer(function (req, res) { pool.connect(async (err, client, done) => { if (err) throw err; const query = new QueryStream(querystring) const stream = client.query(query); stream.on('end', done) return stream.pipe(JSONStream.stringify()).pipe(res) }) }) return server; }
客户端调用代码(src/lib/api.js)
async function getDatabaseData (/** @type {string} */ url) { // fetch internal server api. const response = await fetch(url, { method: 'GET', }); let d = await response.json(); return d.message; }
已知返回server是错误尝试,目前已能在服务端正确获取数据,但无法通过字符串化的Response将数据流式传输到客户端。
解决方案
1. 服务端:返回标准流式Response对象
SvelteKit的RequestHandler要求返回标准Response实例,不能返回http服务器对象。需将PostgreSQL查询流转换为ReadableStream后包装成Response返回:
/** @type {import('./$types').RequestHandler} */ export async function GET({ url }) { return new Response( new ReadableStream({ async start(controller) { const client = await pool.connect(); try { const query = new QueryStream(querystring); const stream = client.query(query); stream.on('data', (chunk) => { // 将单条数据转为JSON字符串,用换行分隔 controller.enqueue(new TextEncoder().encode(JSON.stringify(chunk) + '\n')); }); stream.on('end', () => { controller.close(); client.release(); // 释放数据库连接 }); stream.on('error', (err) => { controller.error(err); client.release(); }); } catch (err) { controller.error(err); client?.release(); } } }), { headers: { 'Content-Type': 'application/x-ndjson', // 换行分隔的JSON格式,适配流式解析 'Transfer-Encoding': 'chunked' } } ); }
2. 客户端:流式解析响应数据
原代码用response.json()会等待完整响应,无法实现流式传输。需逐行解析响应流:
async function getDatabaseData(url) { const response = await fetch(url, { method: 'GET' }); if (!response.body) throw new Error('响应无数据流'); 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); // 按换行分割数据块,处理完整行 const lines = buffer.split('\n'); buffer = lines.pop() || ''; // 保留不完整的最后一行 for (const line of lines) { if (line.trim()) { const data = JSON.parse(line); // 这里添加单条数据的处理逻辑,比如更新UI console.log('收到数据:', data); } } } // 处理最后剩余的完整数据 if (buffer.trim()) { const data = JSON.parse(buffer); console.log('收到最后一条数据:', data); } }
关键说明
- 服务端通过
ReadableStream包装PG查询流,直接返回符合要求的Response对象 - 采用
application/x-ndjson格式,避免一次性传输整个JSON数组,实现真正的流式传输 - 客户端通过
ReadableStreamDefaultReader逐块读取响应,实时解析每一行数据
内容的提问来源于stack exchange,提问作者fm84
相关产品推荐
相关产品推荐

