从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); });
核心逻辑说明
- 追加Blob流式写入:通过
getAppendStream()直接获取云端可写流,将CSV行实时写入Blob,无需本地存储 - 异步流控制:在
data事件中暂停可读流,等待当前行写入完成后再恢复,避免大文件处理时内存溢出 - 分块切换:达到指定记录数后关闭当前流,创建新的追加Blob开始下一块写入
- 收尾处理:在
end事件中关闭最后一个分块的流,确保所有数据都被持久化到云端
内容的提问来源于stack exchange,提问作者Utkarsh
相关产品推荐
相关产品推荐

