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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 03:38:19