Azure虚拟机上Python流式读写Azure Blob存储CSV方案问询
我之前在Azure VM上处理Blob里的大CSV时也碰到过同样的问题,本地靠下载文件的方式完全不适用,后来用流式读写的方案完美解决了——全程不用落地本地文件,直接在内存里完成读取、处理、上传的全流程,给你分享下具体实现:
核心实现思路
利用Azure Blob Storage SDK获取CSV文件的字节流,直接传给pandas读取;处理完成后再把结果转为字节流,上传回Blob存储,全程通过内存流(BytesIO)中转,完全绕开本地文件IO。
具体代码实现
1. 先安装依赖包
pip install azure-storage-blob pandas
2. 完整代码示例
from azure.storage.blob import BlobServiceClient import pandas as pd from io import BytesIO # 替换为你的Blob存储信息 BLOB_CONNECTION_STR = "your_blob_storage_connection_string" CONTAINER_NAME = "your_container_name" SOURCE_BLOB_PATH = "data/source.csv" # 源CSV在Blob中的路径 TARGET_BLOB_PATH = "data/processed.csv" # 处理后要保存的路径 # 初始化Blob服务客户端 blob_service_client = BlobServiceClient.from_connection_string(BLOB_CONNECTION_STR) # ---------------------- 流式读取Blob中的CSV ---------------------- source_blob_client = blob_service_client.get_blob_client( container=CONTAINER_NAME, blob=SOURCE_BLOB_PATH ) # 将Blob内容读取到内存字节流 csv_stream = BytesIO(source_blob_client.download_blob().readall()) # pandas直接读取内存流 df = pd.read_csv(csv_stream) # ---------------------- 你的数据处理逻辑 ---------------------- # 这里替换成你实际的处理代码,示例:计算新增列 df['total'] = df['col1'] + df['col2'] df['processed_at'] = pd.Timestamp.now() # ---------------------- 流式上传处理后的结果 ---------------------- # 将处理后的DataFrame转为CSV字节流 output_stream = BytesIO() df.to_csv(output_stream, index=False, encoding='utf-8') output_stream.seek(0) # 重置流指针到起始位置,否则上传会从末尾开始 # 上传到目标Blob target_blob_client = blob_service_client.get_blob_client( container=CONTAINER_NAME, blob=TARGET_BLOB_PATH ) # overwrite=True 表示覆盖已存在的文件,按需调整 target_blob_client.upload_blob(output_stream, overwrite=True)
关键注意事项
- 大文件优化:如果你的CSV文件特别大(比如几十GB),
readall()会一次性把整个文件加载到内存,可能导致OOM。这种情况下可以用分块读取:# 分块读取Blob并处理 chunk_size = 100000 # 每块10万行 with BytesIO() as chunk_stream: for chunk in source_blob_client.download_blob().chunks(): chunk_stream.write(chunk) chunk_stream.seek(0) df_chunk = pd.read_csv(chunk_stream, chunksize=chunk_size) for sub_chunk in df_chunk: # 处理子块数据 sub_chunk['processed'] = True # 这里可以把处理后的子块追加到目标Blob,需要注意CSV表头只写一次 # 逻辑略复杂,可根据实际需求实现 chunk_stream.seek(0) chunk_stream.truncate() - 权限配置:确保Azure VM具有访问Blob存储的权限,比如给VM分配
Storage Blob Data Contributor角色,或者使用SAS令牌替代连接字符串(更安全)。 - 流指针重置:上传前一定要调用
output_stream.seek(0),否则上传的内容是空的——因为写入DataFrame后,流指针停在末尾。
这样整个流程就完全脱离本地文件系统了,非常适合在VM上运行计算密集型的数据处理任务。
内容的提问来源于stack exchange,提问作者Pepin Peng
相关产品推荐
相关产品推荐

