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

使用Python Azure Functions处理Blob存储CSV文件的实现咨询

Python Azure Functions 实现Blob存储CSV增量处理方案

整体流程通过Blob触发器监听源容器的文件变动,通过入站/出站绑定省略Blob存储手动鉴权、读写的冗余代码,配合偏移量标记实现仅处理新增记录,最终把处理结果写入目标容器。

一、前置准备

  • 在Azure存储账户创建3个容器:source-csv 存放原始CSV文件、target-result 存放处理后的结果文件、offset-meta 存放读取偏移量元数据,用于识别新增未处理记录
  • 本地/云端函数环境安装依赖包:pip install azure-functions pandas
  • 函数应用配置中添加存储账户连接字符串,键名保持为AzureWebJobsStorage(可直接复用函数默认运行存储,也可单独配置业务存储连接串)

二、绑定配置说明

以下实现基于官方推荐的Python v2编程模型,无需单独编写function.json配置文件,绑定直接通过代码装饰器声明。如果使用旧版v1模型,可参考文末的适配配置。

三、完整实现代码

import azure.functions as func
import pandas as pd
import os
import io
import json

# 初始化函数应用实例
app = func.FunctionApp()

# 常量配置,可根据实际环境修改
STORAGE_CONN = os.getenv("AzureWebJobsStorage")
SOURCE_CONTAINER = "source-csv"
TARGET_CONTAINER = "target-result"
META_CONTAINER = "offset-meta"
OFFSET_BLOB = "process_offset.json"

# 声明入站、出站绑定
@app.blob_trigger(
    arg_name="source_blob",
    path=f"{SOURCE_CONTAINER}/{{file_name}}",
    connection="AzureWebJobsStorage"
)
@app.blob_input(
    arg_name="offset_blob",
    path=f"{META_CONTAINER}/{OFFSET_BLOB}",
    connection="AzureWebJobsStorage"
)
@app.blob_output(
    arg_name="target_blob",
    path=f"{TARGET_CONTAINER}/processed_{{file_name}}",
    connection="AzureWebJobsStorage"
)
@app.blob_output(
    arg_name="new_offset_blob",
    path=f"{META_CONTAINER}/{OFFSET_BLOB}",
    connection="AzureWebJobsStorage"
)
def csv_process_flow(
    source_blob: func.InputStream,
    offset_blob: func.InputStream,
    target_blob: func.Out[bytes],
    new_offset_blob: func.Out[bytes]
):
    # 1. 读取历史处理偏移量
    try:
        offset_record = json.loads(offset_blob.read().decode("utf-8"))
        last_processed_pos = offset_record.get(source_blob.name, 0)
    except Exception:
        # 首次运行无偏移量文件时,默认从文件起始位置读取
        last_processed_pos = 0

    # 2. 仅读取未处理的新增内容
    source_blob.seek(last_processed_pos)
    new_content = source_blob.read()
    if not new_content:
        # 无新增内容直接终止流程
        return

    # 3. 自定义数据处理逻辑
    df_new = pd.read_csv(io.BytesIO(new_content))
    # 示例处理:去重、空值替换、字段过滤,可按业务需求修改
    df_processed = df_new.drop_duplicates().fillna("")
    processed_result = df_processed.to_csv(index=False, encoding="utf-8")

    # 4. 写出处理结果到目标容器
    target_blob.set(processed_result.encode("utf-8"))

    # 5. 更新偏移量,标记本次处理完成的位置
    current_pos = last_processed_pos + len(new_content)
    offset_record[source_blob.name] = current_pos
    new_offset_blob.set(json.dumps(offset_record).encode("utf-8"))

四、关键逻辑说明

  • 入站绑定规则:共配置2个入站绑定,第一个是Blob触发器,自动监听源容器下的文件新增、修改事件,触发函数执行;第二个是偏移量文件输入绑定,自动拉取历史处理位置,无需手动编写Blob鉴权、读取代码。
  • 出站绑定规则:共配置2个出站绑定,第一个绑定目标结果容器,直接通过set()方法写入处理后的CSV内容;第二个绑定元数据容器,自动更新处理偏移量,无需手动实现Blob上传逻辑。
  • 增量逻辑适配:如果你的场景是往源容器上传全新CSV文件(而非在原有CSV文件上追加内容),可以把偏移量逻辑调整为记录已处理的文件名列表,直接跳过已处理的文件即可,逻辑更简单。
  • 性能注意事项:如果单CSV文件大小超过100MB,建议把pandas全量读取改为逐块读取,避免函数内存超限;追加写入CSV时要保证新增内容为完整行,避免截断行导致CSV解析失败。

五、v1编程模型适配

如果使用旧版v1编程模型,直接在函数目录下的function.json中写入以下绑定配置即可,上层业务逻辑和上述代码完全一致:

{
  "scriptFile": "__init__.py",
  "bindings": [
    {
      "name": "source_blob",
      "type": "blobTrigger",
      "direction": "in",
      "path": "source-csv/{file_name}",
      "connection": "AzureWebJobsStorage"
    },
    {
      "name": "offset_blob",
      "type": "blob",
      "direction": "in",
      "path": "offset-meta/process_offset.json",
      "connection": "AzureWebJobsStorage"
    },
    {
      "name": "target_blob",
      "type": "blob",
      "direction": "out",
      "path": "target-result/processed_{file_name}",
      "connection": "AzureWebJobsStorage"
    },
    {
      "name": "new_offset_blob",
      "type": "blob",
      "direction": "out",
      "path": "offset-meta/process_offset.json",
      "connection": "AzureWebJobsStorage"
    }
  ]
}

内容的提问来源于stack exchange,提问作者Dataholic

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 16:09:10