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

如何从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 });
    }
}

关键说明

  1. 流的转换:GridFS的下载流和本地文件流都是Node.js Readable类型,直接用原有的ReadableStream封装逻辑即可,无需额外转换。
  2. 响应头设置:通过NextResponse的headers配置替代Pages路由的res.writeHead,直接从GridFS文件元信息中获取contentType和文件大小length。
  3. 错误处理:增加参数校验、文件不存在判断和全局异常捕获,提升路由健壮性。
  4. 断点续传:添加Accept-Ranges: bytes头可支持视频断点续传,优化用户体验。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 00:35:03