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

请求指导:使用Azure Data Factory V2与Python版Azure Functions处理原始文件

Hey there, I’ve got you covered with a step-by-step guide to using Azure Data Factory (ADF) V2 and Python-based Azure Functions for reading, writing, and moving raw files on Azure. Let’s dive right in—no vague explanations, just actionable steps and code you can use:

Prerequisites

First, make sure you have these set up before starting:

  • An active Azure subscription
  • An Azure Data Factory V2 instance already provisioned
  • An Azure Function App with the Python runtime (I recommend the Consumption plan for cost efficiency)
  • An Azure Storage Account (Blob Storage is ideal for raw file storage)
  • The azure-storage-blob package installed in your Function App (add it via requirements.txt)
Step 1: Build Your Python Azure Function

Create an HTTP-triggered function—this will be the endpoint ADF calls to trigger file operations. Below is a complete function that handles all three operations (read, write, move):

import azure.functions as func
from azure.storage.blob import BlobServiceClient
import os

def main(req: func.HttpRequest) -> func.HttpResponse:
    try:
        # Pull parameters sent from ADF in the request body
        req_body = req.get_json()
        storage_conn_str = os.environ["STORAGE_CONNECTION_STRING"]  # Stored in Function App settings
        source_container = req_body.get("source_container")
        source_blob_name = req_body.get("source_blob_name")
        target_container = req_body.get("target_container")
        target_blob_name = req_body.get("target_blob_name")
        operation = req_body.get("operation")  # Accepts "read", "write", "move"

        # Initialize connection to your storage account
        blob_service_client = BlobServiceClient.from_connection_string(storage_conn_str)

        if operation == "read":
            # Fetch and read the source blob
            source_blob_client = blob_service_client.get_blob_client(container=source_container, blob=source_blob_name)
            blob_content = source_blob_client.download_blob().readall()
            # You can add processing logic here (e.g., parse CSV, transform data)
            return func.HttpResponse(f"Successfully read blob. Content length: {len(blob_content)} bytes", status_code=200)

        elif operation == "write":
            # Write content to the target blob (customize this with your actual content)
            target_blob_client = blob_service_client.get_blob_client(container=target_container, blob=target_blob_name)
            # Replace this with content from a read operation, ADF payload, or external source
            content_to_write = b"Raw file content to save to Azure Blob Storage"
            target_blob_client.upload_blob(content_to_write, overwrite=True)
            return func.HttpResponse(f"Successfully wrote to blob: {target_blob_name}", status_code=200)

        elif operation == "move":
            # Move = copy blob to target, then delete the source
            source_blob_client = blob_service_client.get_blob_client(container=source_container, blob=source_blob_name)
            target_blob_client = blob_service_client.get_blob_client(container=target_container, blob=target_blob_name)

            # Start copy operation
            copy_operation = target_blob_client.start_copy_from_url(source_blob_client.url)
            # Wait for copy to complete (add retries/timeouts for production use)
            copy_status = target_blob_client.get_blob_properties().copy.status
            while copy_status != "success":
                copy_status = target_blob_client.get_blob_properties().copy.status

            # Delete the original blob
            source_blob_client.delete_blob()
            return func.HttpResponse(f"Successfully moved blob from {source_blob_name} to {target_blob_name}", status_code=200)

        else:
            return func.HttpResponse(f"Invalid operation: {operation}. Use 'read', 'write', or 'move'", status_code=400)

    except Exception as e:
        return func.HttpResponse(f"Error occurred: {str(e)}", status_code=500)

Quick tips for this function:

  • Store your Storage Account connection string in the Function App’s Configuration (under Application Settings) with the key STORAGE_CONNECTION_STRING—never hardcode secrets!
  • Add azure-storage-blob>=12.0.0 to your requirements.txt so the package installs automatically on deploy.
Step 2: Connect ADF V2 to Your Azure Function

Now let’s link ADF to your function to trigger it:

  1. Create a Web Linked Service:

    • In your ADF instance, go to Manage > Linked Services > New
    • Search for "Web" and select it
    • Configure the settings:
      • URL: Paste your function’s trigger URL (get this from the Function App portal > Your Function > Get Function Url)
      • Authentication: Choose "Function Key"
      • Function Key: Paste the default function key (from Function App portal > Your Function > Manage > Function Keys)
    • Save the linked service (name it something like AzureFunctionLinkedService)
  2. Build a Pipeline to Trigger the Function:

    • Go to Author > Pipelines > New Pipeline
    • Add a Web Activity to the pipeline canvas
    • Configure the activity:
      • Linked Service: Select the AzureFunctionLinkedService you created
      • Method: POST
      • Body: Add a JSON payload with your operation details. For example, to move a file:
        {
            "source_container": "raw-inputs",
            "source_blob_name": "daily-data.csv",
            "target_container": "processed-outputs",
            "target_blob_name": "daily-data-20240520.csv",
            "operation": "move"
        }
        
    • Save the pipeline and run a debug test to verify it works. Check the Function App’s logs (Monitor > Logs) and ADF’s pipeline run details for any issues.
Step 3: Customize for Your Workflow
  • Reading files: Modify the "read" block to process the blob content (e.g., parse JSON/CSV, clean data) instead of just returning the length.
  • Writing files: Replace the hardcoded content_to_write with data from ADF’s payload, a database query, or another blob.
  • Moving files: For large files, add robust error handling (timeouts, retries) instead of the simple while loop to avoid hanging.
Troubleshooting Common Issues
  • Permission errors: Assign the Storage Blob Data Contributor role to your Function App’s managed identity in the Storage Account’s Access Control (IAM) section.
  • Function not triggering: Double-check the Web Activity’s URL and function key—typos are a common culprit.
  • Copy operation failures: Ensure the source blob exists and the target container is created. Check the blob’s properties in the Storage Account to view copy status details.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:04:21