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

如何直接从Service Bus消息启动Azure Durable Functions编排?

直接通过Service Bus触发Azure Durable Functions编排函数的实现方案

可行性确认

完全可以直接通过Service Bus消息触发编排函数,无需中间转发函数。Azure Durable Functions支持将编排函数与Service Bus触发器(队列/主题)直接绑定,只需调整装饰器配置并处理消息输入即可。

具体实现步骤

  1. 移除中间转发函数:删除原有的process_service_bus_message函数,不再需要它作为中转。
  2. 合并触发器绑定:给编排函数process_workflow同时添加Service Bus主题触发器和编排触发器的装饰器。
  3. 处理消息输入:在编排函数中直接解析Service Bus消息的内容,作为编排的输入数据。

修改后的完整代码示例

import azure.functions as func
import azure.durable_functions as df
import logging
import json
from item_processor import handle_item
from summary_builder import build_summary

# 初始化Durable Functions应用
app_instance = df.DurableApp(http_auth_level=func.AuthLevel.ANONYMOUS)

# 给编排函数同时绑定Service Bus主题触发器和编排触发器
@app_instance.service_bus_topic_trigger(
    arg_name="incoming_data",
    topic_name="item-topic",
    subscription_name="main-subscription",
    connection="ServiceBusConnectionString"
)
@app_instance.orchestration_trigger(context_name="workflow_context")
def process_workflow(workflow_context: df.DurableOrchestrationContext, incoming_data: func.ServiceBusMessage):
    # 解析Service Bus消息体
    data_payload = incoming_data.get_body().decode('utf-8')
    logging.info(f"Received Service Bus message: {data_payload}")
    data_dict = json.loads(data_payload)
    
    # 原编排逻辑保持不变
    # 并行处理所有元素
    task_list = [workflow_context.call_activity("handle_item", item) for item in data_dict['elements']]
    processed_results = yield workflow_context.task_all(task_list)

    # 聚合结果
    results_aggregation = {
        "process_id": data_dict['processId'],
        "results": processed_results
    }

    # 生成汇总报告
    final_summary = yield workflow_context.call_activity("build_summary", results_aggregation)
    logging.info(f"Orchestration completed. Summary: {final_summary}")

    return final_summary

# 活动函数配置保持不变
handle_item = app_instance.activity_trigger(input_name="item")(handle_item)
build_summary = app_instance.activity_trigger(input_name="summary")(build_summary)

关键注意事项

  • 编排函数的确定性:务必确保编排函数内只执行确定性操作(如调用活动函数、使用上下文提供的API),不要在编排函数中直接调用外部服务、生成随机数或获取当前时间(这些操作要放到活动函数里),否则会导致编排重放时出现不一致问题。
  • 消息序列化:Service Bus消息的内容必须是可JSON序列化的,因为编排的输入会被持久化到Durable Task的存储账户中。
  • 权限配置:确保函数应用的身份(如系统分配的托管身份)拥有Service Bus主题的Listen权限,以及Durable Functions所需的存储账户权限。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 14:37:15