如何使用Durable Functions向Blob服务的CSV文件追加写入数据
解决方案:Azure Durable Functions Python v2 追加CSV到Blob存储
要实现CSV数据追加,不能直接用outputblob.set()(该方法会覆盖原有文件),得通过Blob Storage SDK手动处理「读取现有内容+合并新数据+写回Blob」的逻辑,或者利用追加Blob类型优化操作。以下是具体实现步骤:
1. 配置依赖
在requirements.txt中添加必要包:
azure-functions>=1.17.0 pandas>=2.0.0 azure-storage-blob>=12.18.0
2. 实现活动函数(核心追加逻辑)
写一个Durable活动函数,负责处理Blob的读取、数据合并和写入:
import azure.functions as func import pandas as pd from azure.storage.blob import BlobServiceClient import os from io import StringIO @app.activity_trigger(input_name="input") def append_csv_to_blob(input: dict) -> str: # 从输入参数获取配置和数据 blob_container = input["container_name"] blob_name = input["blob_name"] df_data = input["df_data"] # 初始化Blob客户端(从环境变量取连接字符串) connect_str = os.getenv("AzureWebJobsStorage") blob_service_client = BlobServiceClient.from_connection_string(connect_str) blob_client = blob_service_client.get_blob_client(container=blob_container, blob=blob_name) # 构造待追加的DataFrame df = pd.DataFrame(df_data) csv_buffer = StringIO() if blob_client.exists(): # Blob已存在:读取原有内容,新数据仅写内容(跳过表头) existing_csv = blob_client.download_blob().readall().decode("utf-8") df.to_csv(csv_buffer, index=False, header=False) combined_csv = existing_csv + "\n" + csv_buffer.getvalue().strip() else: # Blob不存在:直接写入完整CSV(包含表头) df.to_csv(csv_buffer, index=False) combined_csv = csv_buffer.getvalue() # 写回Blob(覆盖原文件,因为已处理追加逻辑) blob_client.upload_blob(combined_csv, overwrite=True) return f"数据已成功追加到Blob: {blob_name}"
3. 编写编排函数调用活动函数
import azure.durable_functions as df @app.orchestration_trigger(context_name="context") def orchestrator_function(context: df.DurableOrchestrationContext): # 定义Blob信息和待追加的DataFrame数据 task_input = { "container_name": "你的容器名", "blob_name": "目标文件.csv", "df_data": { "列1": [4,5,6], "列2": ["d","e","f"] } } # 调用活动函数执行追加 result = yield context.call_activity("append_csv_to_blob", task_input) return result
关键注意事项
- 避免重复表头:Blob已存在时,新DataFrame导出CSV必须跳过表头,否则文件会重复出现表头行。
- 连接字符串配置:确保本地
local.settings.json或云端函数应用的应用设置中,AzureWebJobsStorage已配置为你的Blob存储连接字符串。 - 高效追加方案:如果你的Blob是追加Blob(创建时指定类型为AppendBlob),可以直接用
blob_client.append_block()写入,无需读取整个文件,适合高频追加场景:
# 追加Blob模式示例 if not blob_client.exists(): blob_client.create_append_blob() # 首次写入带表头 df.to_csv(csv_buffer, index=False) blob_client.append_block(csv_buffer.getvalue()) else: # 后续仅追加数据(无表头) df.to_csv(csv_buffer, index=False, header=False) blob_client.append_block(csv_buffer.getvalue())
内容的提问来源于stack exchange,提问作者TY00
相关产品推荐
相关产品推荐

