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

如何将Blob触发的Starter Function输入传递给Azure Durable Functions编排器活动函数?

解决Durable Functions中传递Blob InputStream的问题

直接传递func.InputStream对象给Durable Orchestrator会因为无法JSON序列化报错,正确的做法是传递Blob的定位信息,再在活动函数中重新获取Blob内容。

步骤1:修改Starter函数,传递可序列化的Blob元数据

把client_input改成包含Blob路径和连接字符串配置名的字典,不要直接传InputStream:

async def main(myblob: func.InputStream, starter: str) -> func.InputStream:
    logging.info(f"Python blob trigger function processed blob \n"
                 f"Name: {myblob.name}\n"
                 f"Blob Size: {myblob.length} bytes\n\n")        
    client = df.DurableOrchestrationClient(starter)
    # 传递Blob的路径和连接字符串配置键(不要直接传连接字符串明文)
    blob_info = {
        "blob_path": myblob.name,
        "connection_name": "AzureWebJobsStorage"  # 对应functions.json里的连接配置名
    }
    instance_id = await client.start_new('Orchestrator', client_input=blob_info)
    logging.info(f"Started orchestration with ID = '{instance_id}'.")
    return myblob

步骤2:编排器函数传递元数据给活动函数

编排器函数直接把收到的Blob元数据转发给活动函数即可:

def orchestrator_function(context: df.DurableOrchestrationContext):
    blob_info = context.get_input()
    # 调用活动函数,传入Blob元数据
    result = yield context.call_activity("ProcessBlobActivity", blob_info)
    return result

步骤3:活动函数获取Blob内容(两种方式)

方式1:使用Blob输入绑定

修改活动函数的functions.json,通过绑定直接获取InputStream:

{
  "scriptFile": "__init__.py",
  "bindings": [
    {
      "name": "myblob",
      "type": "blob",
      "direction": "in",
      "path": "{blob_path}",
      "connection": "{connection_name}"
    },
    {
      "name": "activityTrigger",
      "type": "activityTrigger",
      "direction": "in"
    }
  ]
}

活动函数代码:

def main(activity_input: dict, myblob: func.InputStream):
    logging.info(f"Processing blob: {myblob.name}")
    # 这里可以直接把myblob传给认知服务,比如如果服务接受流:
    # cognitive_service_client.process_stream(myblob)
    return "Blob processed successfully"

方式2:使用Azure Storage SDK手动读取

如果需要更灵活的控制(比如读取部分内容、转换格式),可以用SDK读取Blob:

from azure.storage.blob import BlobServiceClient
import os

def main(activity_input: dict):
    blob_path = activity_input["blob_path"]
    connection_name = activity_input["connection_name"]
    # 从环境变量获取连接字符串
    connection_string = os.environ[connection_name]
    
    blob_service_client = BlobServiceClient.from_connection_string(connection_string)
    # 拆分Blob路径为容器和文件名
    container_name, blob_name = blob_path.split('/', 1)
    blob_client = blob_service_client.get_blob_client(container=container_name, blob=blob_name)
    
    # 获取Blob流(可以直接传给认知服务)
    with blob_client.open_read() as blob_stream:
        # cognitive_service_client.process_stream(blob_stream)
        logging.info(f"Processed blob stream: {blob_name}")
    
    return "Blob processed successfully"

关键注意事项

  • 不要直接传递连接字符串明文,用配置名从环境变量读取,保证安全性。
  • 两种方式都能获取到可被认知服务使用的Blob流/字节,不会丢失功能。
  • 错误TypeError: class <class 'azure.functions.blob.InputStream'> does not expose a 'to_json' function的根源就是InputStream无法序列化,传递元数据是标准解决方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 03:35:27