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

如何在Next.js App Router中实现文件内容流式传输至客户端?

Next.js App Router 实现SSH日志流式传输

问题背景

我正在使用Next.js App Router的路由处理程序向客户端传输文件内容,目前采用轮询方式刷新日志,希望改用流式传输方案优化体验。现有代码通过ssh2执行tail -n获取指定行数的静态日志,但不知道如何适配tail -f实现实时流式输出,也不清楚如何将回调式的数据流转换为异步生成器,同时对客户端如何消费该流存在疑问。

现有后端代码:

// app/api/route.ts
async function getLogs(lines: number): Promise<string> {
  const conn = new ssh2.Client();

  return new Promise((resolve, reject) => {
    conn.connect(sshConfig);
    conn.on("ready", function () {
      console.log("Client :: ready");
      conn.exec(
        `tail -${lines} ${filePath}`,
        function (err, stream) {
          if (err) throw err;
          stream
            .on("close", function (code: string, signal: string) {
              conn.end();
            })
            .on("data", function (data: string) {
              resolve(data);
            })
            .stderr.on("data", function (data) {
              reject(data);
            });
        }
      );
    });
  });
}

export async function GET(request: Request) {
  const { searchParams } = new URL(request.url);
  const lines = searchParams.get("lines") ?? "10";
  const data = await getLogs(Number(lines));

  return Response.json({ logs: data.toString() });
}

解决方案

一、后端改造:实现流式日志输出

要实现tail -f的实时流式传输,核心是将回调式的SSH数据流转换为异步生成器,再包装成ReadableStream返回给客户端。

1. 改造日志获取函数为异步生成器

将原有的getLogs替换为返回异步生成器的streamLogs,通过队列缓存数据流,解决回调与生成器的适配问题:

// app/api/route.ts
import { Client } from 'ssh2';

async function* streamLogs(lines: number) {
  const conn = new Client();
  const encoder = new TextEncoder();

  // 初始化SSH连接
  await new Promise<void>((resolve, reject) => {
    conn.connect(sshConfig);
    conn.on('ready', resolve);
    conn.on('error', reject);
  });

  console.log('Client :: ready');

  // 执行tail命令:先返回指定行数历史日志,再实时监听新增内容
  const stream = await new Promise<any>((resolve, reject) => {
    conn.exec(`tail -n ${lines} -f ${filePath}`, (err, stream) => {
      if (err) {
        conn.end();
        reject(err);
        return;
      }
      resolve(stream);
    });
  });

  const dataQueue: Buffer[] = [];
  let isClosed = false;

  // 监听日志数据流,存入队列
  stream.on('data', (data: Buffer) => {
    dataQueue.push(data);
  });

  // 流关闭时标记状态并断开SSH连接
  stream.on('close', () => {
    isClosed = true;
    conn.end();
  });

  // 处理命令执行错误
  stream.stderr.on('data', (errData: Buffer) => {
    throw new Error(errData.toString());
  });

  // 循环读取队列数据并输出
  while (!isClosed || dataQueue.length > 0) {
    if (dataQueue.length > 0) {
      const data = dataQueue.shift();
      if (data) {
        yield encoder.encode(data.toString());
      }
    } else {
      // 短时间等待避免空循环占用CPU
      await new Promise(resolve => setTimeout(resolve, 100));
    }
  }
}

2. 路由处理函数返回流

复用工具函数将异步生成器转换为ReadableStream并返回:

// 工具函数:将异步生成器转换为ReadableStream
function iteratorToStream(iterator: AsyncGenerator<Uint8Array>) {
  return new ReadableStream({
    async pull(controller) {
      try {
        const { value, done } = await iterator.next();
        if (done) {
          controller.close();
        } else {
          controller.enqueue(value);
        }
      } catch (err) {
        controller.error(err);
      }
    },
  });
}

export async function GET(request: Request) {
  const { searchParams } = new URL(request.url);
  const lines = searchParams.get('lines') ?? '10';

  try {
    const logIterator = streamLogs(Number(lines));
    const stream = iteratorToStream(logIterator);

    return new Response(stream, {
      headers: {
        'Content-Type': 'text/plain; charset=utf-8',
        'Transfer-Encoding': 'chunked',
      },
    });
  } catch (err) {
    return new Response(JSON.stringify({ error: (err as Error).message }), {
      status: 500,
      headers: { 'Content-Type': 'application/json' },
    });
  }
}

二、客户端消费流式日志

客户端通过fetch获取流,使用ReadableStream.getReader()逐块读取数据,实时更新UI。

React组件示例(客户端组件)

'use client';

import { useEffect, useState } from 'react';

export default function LogViewer() {
  const [logs, setLogs] = useState('');
  const [error, setError] = useState('');

  useEffect(() => {
    let reader: ReadableStreamDefaultReader | null = null;

    async function fetchLogs() {
      try {
        const response = await fetch('/api/logs?lines=10');
        if (!response.ok) {
          throw new Error(`请求失败:${response.status}`);
        }

        const stream = response.body;
        if (!stream) {
          throw new Error('服务器未返回数据流');
        }

        reader = stream.getReader();
        const decoder = new TextDecoder('utf-8');

        // 持续读取流数据
        while (true) {
          const { done, value } = await reader.read();
          if (done) break;

          // 追加新日志到现有内容
          setLogs(prev => prev + decoder.decode(value));
        }
      } catch (err) {
        setError((err as Error).message);
      } finally {
        if (reader) {
          reader.releaseLock();
        }
      }
    }

    fetchLogs();

    // 组件卸载时取消流读取,避免内存泄漏
    return () => {
      if (reader) {
        reader.cancel();
      }
    };
  }, []);

  return (
    <div className="p-4">
      {error && <div className="text-red-500 mb-4">{error}</div>}
      <pre className="bg-gray-100 p-4 rounded overflow-auto max-h-[400px] whitespace-pre-wrap">
        {logs}
      </pre>
    </div>
  );
}

关键说明

  • 后端使用tail -n ${lines} -f组合命令,先返回指定行数的历史日志,再实时输出新增内容
  • 异步生成器通过队列缓存数据流,解决回调函数与生成器的适配问题
  • 客户端通过ReadableStream.getReader()逐块读取流数据,实现日志实时更新
  • 组件卸载时需主动取消流读取,避免内存泄漏

内容的提问来源于stack exchange,提问作者Reifocs

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 22:24:58