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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 05:32:43