You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在服务端用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 19:42:34