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

从Azure Blob存储流式读取CSV并分块直传至Blob存储的方案咨询

问题描述

我需要从Azure Blob存储通过可读流(readableStream)读取包含百万条记录的CSV文件,将其按每10k/20k条记录分块,再上传回Azure Blob存储作为独立文件。当前采用的方案是流式读取数据后在本地生成分块文件,再逐一上传至存储账户。请问是否存在直接在Blob存储对象中创建并写入分块的方法?

我已尝试以下代码:

import csvParser from "csv-parser";

let newchunk: boolean = true, currentChunk: number =0, currentIndex: number = 1, chunkWriteStream: any;
const writeStream = azureBlobObject.download(0).readableStreamBody; 
writeStream.pipe(csvParser())
                .on('data', (data: any) => {

                    if(newchunk) {
                        newchunk = false;
                        const newContainerClient = this.AzureBlobConInstance.getContainerClient(container);
                        const newBlobClient: BlockBlobClient = newContainerClient.getBlockBlobClient(`file${currentChunk}`);
                        chunkWriteStream = (newBlobClient.getAppendBlobClient() as any).getBlobAppendStream()//fs.createWriteStream(name, { flags: "a" });
                        chunkWriteStream.write(`${Object.keys(data).join(",")}\n`);
                    }

                    chunkWriteStream.write(`${Object.values(data).join(",")}\n`);

                    if (currentIndex >= chunkSize) {
                        newchunk = true;
                        currentChunk++;
                        currentIndex = 1;
                    } else {
                        currentIndex++;
                    }
                })

解决方案

可以直接通过Azure Blob的**追加Blob(Append Blob)**流式写入能力实现,无需本地文件中转。你的思路方向正确,但需要修正流处理、收尾逻辑等问题,以下是优化后的实现:

关键修正点

  • 追加Blob首次写入前必须调用createIfNotExists初始化
  • 流写入是异步操作,需用Promise包裹避免数据积压或丢失
  • 必须处理csv-parser的end事件,确保最后一个分块的流被正确关闭
  • 需暂停/恢复可读流,防止内存溢出

改进后的代码实现

import csvParser from "csv-parser";
import { AppendBlobClient, ContainerClient } from "@azure/storage-blob";

// 配置分块大小
const CHUNK_SIZE = 10000; // 每块10k条记录
let currentChunk = 0;
let currentIndex = 1;
let currentAppendBlob: AppendBlobClient | null = null;
let currentWriteStream: any = null;

async function writeToChunk(data: any) {
    // 初始化新分块
    if (!currentAppendBlob || currentIndex === 1) {
        const containerClient: ContainerClient = this.AzureBlobConInstance.getContainerClient(container);
        currentAppendBlob = containerClient.getAppendBlobClient(`file${currentChunk}`);
        // 初始化追加Blob(首次写入必须执行)
        await currentAppendBlob.createIfNotExists();
        currentWriteStream = currentAppendBlob.getAppendStream();

        // 写入表头(仅新分块第一次写入时执行)
        if (currentIndex === 1) {
            const header = `${Object.keys(data).join(",")}\n`;
            await new Promise((resolve, reject) => {
                currentWriteStream.write(header, (err: Error | null) => err ? reject(err) : resolve(null));
            });
        }
    }

    // 写入当前数据行
    const row = `${Object.values(data).join(",")}\n`;
    await new Promise((resolve, reject) => {
        currentWriteStream.write(row, (err: Error | null) => err ? reject(err) : resolve(null));
    });

    // 达到分块大小,切换到下一个分块
    if (currentIndex >= CHUNK_SIZE) {
        await new Promise((resolve, reject) => {
            currentWriteStream.end((err: Error | null) => err ? reject(err) : resolve(null));
        });
        currentChunk++;
        currentIndex = 1;
        currentAppendBlob = null;
        currentWriteStream = null;
    } else {
        currentIndex++;
    }
}

// 启动处理流程
const downloadStream = azureBlobObject.download(0).readableStreamBody;
if (!downloadStream) throw new Error("无法获取Blob可读流");

downloadStream
    .pipe(csvParser())
    .on('data', async (data: any) => {
        // 暂停可读流,避免数据积压
        downloadStream.pause();
        try {
            await writeToChunk(data);
        } catch (err) {
            console.error("分块写入失败:", err);
            throw err;
        } finally {
            // 恢复可读流
            downloadStream.resume();
        }
    })
    .on('end', async () => {
        // 处理最后一个未完成的分块
        if (currentWriteStream) {
            await new Promise((resolve, reject) => {
                currentWriteStream.end((err: Error | null) => err ? reject(err) : resolve(null));
            });
            console.log("所有分块处理完成");
        }
    })
    .on('error', (err: Error) => {
        console.error("CSV解析或流处理失败:", err);
    });

核心逻辑说明

  1. 追加Blob流式写入:通过getAppendStream()直接获取云端可写流,将CSV行实时写入Blob,无需本地存储
  2. 异步流控制:在data事件中暂停可读流,等待当前行写入完成后再恢复,避免大文件处理时内存溢出
  3. 分块切换:达到指定记录数后关闭当前流,创建新的追加Blob开始下一块写入
  4. 收尾处理:在end事件中关闭最后一个分块的流,确保所有数据都被持久化到云端

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 13:18:08