如何将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
相关产品推荐
相关产品推荐

