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

Node.js中promisified pipeline处理损坏文件时应用崩溃问题求助

问题修复:处理损坏文件时Node.js上传流程崩溃的解决方案

问题根源

你使用的file.stream()是Web Streams API的ReadableStream,并非Node.js原生的Stream接口。直接将其传入Node.js的pipeline(无论是promisify版本还是原生Promise版),会导致错误无法被try/catch捕获——Web Stream的错误会绕过Node.js的Promise捕获机制,最终触发未捕获的异常导致应用崩溃。

修复步骤及代码实现

1. 核心修复:转换Web Stream为Node.js可读流

使用Node.js 17.0.0+提供的stream.Readable.fromWeb()方法,将Web Stream转换为Node.js兼容的可读流,确保pipeline能正确处理错误传递。

2. 兜底处理:监听流的错误事件+全局异常捕获

额外绑定流的error事件作为兜底,同时添加全局异常捕获,避免极端情况下的进程崩溃。

3. 优化细节:异步文件操作+临时文件清理

用异步fs.mkdir替代同步操作,避免阻塞事件循环;上传失败时自动清理临时文件,防止垃圾文件堆积。

修改后的完整代码:

import { pipeline, Readable } from 'stream';
import { promisify } from 'util';
import fs from 'fs/promises';
import fsSync from 'fs';

const streamPipeline = promisify(pipeline);

try {
    const formData = await req.formData();
    let file;

    if (!formData.has('file')) {
        throw new Error('No file uploaded');
    } else {
        file = formData.get('file');
    }

    // 验证文件类型
    if (!file || typeof file !== 'object' || !file.name || !file.name.endsWith('.avro')) {
        throw new Error('Invalid file type. Only .avro files are allowed');
    }

    const tempDir = './temp';
    await fs.mkdir(tempDir, { recursive: true });

    const tempPath = `${tempDir}/${file.name}`;
    // 将Web Stream转换为Node.js可读流
    const nodeReadableStream = Readable.fromWeb(file.stream());
    const writeStream = fsSync.createWriteStream(tempPath);

    // 显式监听写入流错误,兜底处理
    writeStream.on('error', (err) => {
        logger.error('Write stream error:', err);
        fs.unlink(tempPath).catch(() => {});
    });

    await streamPipeline(nodeReadableStream, writeStream);
} catch (error) {
    logger.error('Upload error:', error);
    // 清理临时文件
    if (typeof tempPath !== 'undefined') {
        fs.unlink(tempPath).catch(() => {});
    }
    return new Response("upload error", {
        status: 500,
    });
}

// 全局兜底:捕获未处理的Promise拒绝和未捕获异常
process.on('unhandledRejection', (reason, promise) => {
    logger.error('Unhandled Rejection at:', promise, 'reason:', reason);
});

process.on('uncaughtException', (err) => {
    logger.error('Uncaught Exception thrown:', err);
    // 可选:执行优雅退出逻辑后再终止进程
    // server.close(() => process.exit(1));
});

低版本Node.js兼容方案

如果你的Node.js版本低于17.0.0,可使用第三方包web-streams-node转换流:

npm install web-streams-node

替换转换逻辑:

import { toNodeReadable } from 'web-streams-node';
// ...
const nodeReadableStream = toNodeReadable(file.stream());

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 09:49:52