如何在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
相关产品推荐
相关产品推荐

