如何从GridFSBucket获取ReadableStream?(Next.js路由迁移)
在Next.js App路由中用GridFSBucket实现视频流式传输
需要将旧Pages路由中基于GridFSBucket的视频流式输出逻辑,迁移到App路由的streamFile函数中,替代原本读取本地文件的fs.createReadStream。不清楚如何将GridFS的下载流转换为App路由所需的ReadableStream<Uint8Array>,或者是否需要额外转换。
原App路由示例代码(读取本地文件)
import { NextRequest, NextResponse } from "next/server"; import { ReadableOptions } from "stream"; function streamFile(path: string, options?: ReadableOptions): ReadableStream<Uint8Array> { const downloadStream = fs.createReadStream(path, options); // 需要替换这部分 return new ReadableStream({ start(controller) { downloadStream.on("data", (chunk: Buffer) => controller.enqueue(new Uint8Array(chunk))); downloadStream.on("end", () => controller.close()); downloadStream.on("error", (error: NodeJS.ErrnoException) => controller.error(error)); }, cancel() { downloadStream.destroy(); }, }); } export async function GET(req: NextRequest): Promise<NextResponse> { const file = req.nextUrl.searchParams.get("file"); const stats: Stats = await fs.promises.stat(file); const data: ReadableStream<Uint8Array> = streamFile(file); const res = new NextResponse(data, { status: 200, headers: new Headers({ "content-type": "video/mp4", "content-length": stats.size + "", }), }); return res; }
旧Pages路由流式输出代码(基于GridFSBucket)
export const getServerSideProps: GetServerSideProps = async ({ res, query: { hash } }) => { const database = await mongodb() const Videos = database.collection('videos') const { fileId } = await Videos.findOne({ uid: hash }) const Files = new GridFSBucket(database) const id = new ObjectId(fileId) const file: GridFSFile = await new Promise((resolve, reject) => { Files.find({ _id: id }).toArray((err, files) => { if (err) reject(err) resolve(files[0]) }) }) const { contentType } = file || {} res.writeHead('Content-Type', contentType) // contentType为"video/mp4" Files.openDownloadStream(id) .on('data', (chunk) => { res.write(chunk) }) .on('end', () => { res.end() }) .on('error', (err) => { throw err }) return { props: {} } }
解决方案
GridFSBucket的openDownloadStream返回的是Node.js原生Readable流,和fs.createReadStream的类型一致,因此可以直接复用原streamFile的封装逻辑,只需要替换流的来源。同时在App路由中不需要直接操作res.writeHead,而是通过NextResponse的headers配置来设置响应头。
修改后的完整代码如下:
import { NextRequest, NextResponse } from "next/server"; import { Readable } from "stream"; import { mongodb } from "@/lib/mongodb"; // 替换为你的MongoDB连接模块路径 import { ObjectId, GridFSBucket, GridFSFile } from "mongodb"; // 改造streamFile,直接接收Node.js Readable流 function streamFile(downloadStream: Readable): ReadableStream<Uint8Array> { return new ReadableStream({ start(controller) { downloadStream.on("data", (chunk: Buffer) => controller.enqueue(new Uint8Array(chunk))); downloadStream.on("end", () => controller.close()); downloadStream.on("error", (error) => controller.error(error)); }, cancel() { downloadStream.destroy(); }, }); } export async function GET(req: NextRequest): Promise<NextResponse> { try { const hash = req.nextUrl.searchParams.get("hash"); if (!hash) { return new NextResponse("缺少hash参数", { status: 400 }); } const database = await mongodb(); const Videos = database.collection("videos"); // 查询对应视频的fileId const videoDoc = await Videos.findOne({ uid: hash }); if (!videoDoc?.fileId) { return new NextResponse("视频不存在", { status: 404 }); } const fileId = new ObjectId(videoDoc.fileId); const bucket = new GridFSBucket(database); // 获取文件元信息(用于Content-Type和Content-Length) const file: GridFSFile | undefined = await new Promise((resolve, reject) => { bucket.find({ _id: fileId }).toArray((err, files) => { if (err) reject(err); resolve(files[0]); }); }); if (!file) { return new NextResponse("文件不存在", { status: 404 }); } // 创建GridFS下载流并转换为Web ReadableStream const downloadStream = bucket.openDownloadStream(fileId); const dataStream = streamFile(downloadStream); return new NextResponse(dataStream, { status: 200, headers: new Headers({ "Content-Type": file.contentType || "video/mp4", "Content-Length": file.length.toString(), "Accept-Ranges": "bytes", // 可选:支持断点续传 }), }); } catch (error) { console.error("流式传输错误:", error); return new NextResponse("服务器内部错误", { status: 500 }); } }
关键说明
- 流的转换:GridFS的下载流和本地文件流都是Node.js Readable类型,直接用原有的
ReadableStream封装逻辑即可,无需额外转换。 - 响应头设置:通过
NextResponse的headers配置替代Pages路由的res.writeHead,直接从GridFS文件元信息中获取contentType和文件大小length。 - 错误处理:增加参数校验、文件不存在判断和全局异常捕获,提升路由健壮性。
- 断点续传:添加
Accept-Ranges: bytes头可支持视频断点续传,优化用户体验。
内容的提问来源于stack exchange,提问作者user3142695
相关产品推荐
相关产品推荐

