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

Azure Durable Orchestrator提前标记完成,未触发全部Activity函数

问题:Durable Function Orchestrator提前完成,未执行全部Activity调用

我开发了一个每日运行的Durable Function,会根据数据库行数调用0-1000次Activity函数。起初运行正常,但近日发现Orchestrator启动几分钟后就被标记为“已完成”,即便还有数百次Activity函数需要调用。当前Orchestrator中的for循环未按预期执行:没有根据DataFrame行数调用DownloadImages-v2_5 Activity函数。今日函数4点启动,4:06就已完成,Orchestrator日志显示“Products left to process: 310”,数据库也确认有310行待处理。


触发Orchestrator的Timer函数

import logging

import azure.functions as func
import azure.durable_functions as df


async def main(mytimer: func.TimerRequest, starter: str) -> None:
    client = df.DurableOrchestrationClient(starter)
    instance_id = await client.start_new("DownloadsOrchestrator-v2_5", None, None)
    logging.info(f"Started orchestration with ID = '{instance_id}'.")

Timer函数的function.json

{
  "scriptFile": "__init__.py",
  "bindings": [
    {
      "name": "mytimer",
      "type": "timerTrigger",
      "direction": "in",
      "schedule": "0 0 4 * * *"
    },
    {
      "name": "starter",
      "type": "durableClient",
      "direction": "in"
    }
  ]
}

Orchestrator函数

import logging
from os import getenv

import azure.functions as func
import azure.durable_functions as df

from pyodbc import connect
from pandas import DataFrame


def orchestrator_function(context: df.DurableOrchestrationContext):

    # Dados para BD:
    server = function-to-get-the-value
    database = function-to-get-the-value
    username = function-to-get-the-value
    password = function-to-get-the-value
    driver = 'ODBC Driver 17 for SQL Server'


    con_str='Driver={ODBC Driver 17 for SQL Server};Server=tcp:'+server+',1433;Database='+database+';Uid='+username+';Pwd='+password+';Encrypt=yes;TrustServerCertificate=no;Connection Timeout=30;Authentication=ActiveDirectoryPassword'

    # SQL que nos dá produtos criados nas últiumas 24h:
    sql = """
    SELECT col1, col2
    FROM table_h
    WHERE date_cre BETWEEN DATEADD(day, -1, SYSDATETIME()) AND SYSDATETIME()
    """

    conn = connect(con_str)

    cursor = conn.cursor()

    cursor.execute(sql)

    rows = cursor.fetchall()

    products_orc = DataFrame.from_records(rows, columns=[col[0] for col in cursor.description])
    
    conn.commit()
    cursor.close()


    prods_loop_orc = products_orc[['Referencia', 'Marca']]
    prods_loop_orc = prods_loop_orc.drop_duplicates(keep="last", inplace=False)
    prods_loop_orc.reset_index(drop=True, inplace=True)

    logging.info("Products left to process: " + str(len(prods_loop_orc.index)))


    # list to return:
    results_list = [] 

    for i in range(0, len(prods_loop_orc)):
        #i=0

        data = products_orc.loc[products_orc['Referencia'] == prods_loop_orc['Referencia'][i]]

        data = data.to_dict()

        result = yield context.call_activity('DownloadImages-v2_5', data)

        results_list.append(result)

    return results_list

main = df.Orchestrator.create(orchestrator_function)

Orchestrator的function.json

{
  "scriptFile": "__init__.py",
  "bindings": [
    {
      "name": "context",
      "type": "orchestrationTrigger",
      "direction": "in"
    }
  ]
}

DownloadImages-v2_5 Activity函数

# 导入多个包并将图片下载到指定文件夹的Python代码

Activity函数的function.json

{
  "scriptFile": "__init__.py",
  "bindings": [
    {
      "name": "name",
      "type": "activityTrigger",
      "direction": "in"
    }
  ]
}

问题原因及解决方法

1. 违反Durable Orchestrator的确定性要求

Durable Orchestrator必须是确定性的,不能在内部直接执行外部服务调用(比如你的数据库连接查询)、使用非确定性API或修改外部状态。直接在Orchestrator里操作数据库会导致重放时出现状态不一致,触发异常终止或提前完成。

修复方案:把数据库查询逻辑移到单独的Activity函数中,让Orchestrator调用该Activity获取数据后再执行循环。

2. 未捕获异常导致循环中断

如果循环中调用Activity或数据转换时发生未捕获异常,Orchestrator可能静默终止并标记为“完成”。比如data.to_dict()转换失败、Activity函数抛出未处理异常等。

修复方案:

  • 在Activity函数中添加完整的日志和异常捕获
  • 在Orchestrator循环中增加try-except块,捕获调用Activity时的异常,避免整个流程终止
  • 检查prods_loop_orc数据是否存在空值、格式错误等问题

3. 执行超时(次要可能)

消费计划下Orchestrator默认超时为10分钟,若数据库查询耗时过长可能触发超时,但你的场景是6分钟完成,更可能是确定性问题导致。

修复方案:将耗时外部操作移到Activity,单独配置Activity的超时时间。


修复后的核心代码示例

重构后的Orchestrator函数

import logging
from os import getenv

import azure.functions as func
import azure.durable_functions as df


def orchestrator_function(context: df.DurableOrchestrationContext):
    # 调用Activity获取数据库数据
    products_orc = yield context.call_activity('FetchProductsFromDB', None)
    
    prods_loop_orc = products_orc[['Referencia', 'Marca']]
    prods_loop_orc = prods_loop_orc.drop_duplicates(keep="last", inplace=False)
    prods_loop_orc.reset_index(drop=True, inplace=True)

    logging.info("Products left to process: " + str(len(prods_loop_orc.index)))

    results_list = [] 

    for i in range(0, len(prods_loop_orc)):
        try:
            data = products_orc.loc[products_orc['Referencia'] == prods_loop_orc['Referencia'][i]]
            data = data.to_dict()
            result = yield context.call_activity('DownloadImages-v2_5', data)
            results_list.append(result)
        except Exception as e:
            logging.error(f"Failed to process product {prods_loop_orc['Referencia'][i]}: {str(e)}")
            continue

    return results_list

main = df.Orchestrator.create(orchestrator_function)

新增的FetchProductsFromDB Activity函数

import logging
from os import getenv
from pyodbc import connect
from pandas import DataFrame

def main(name: str) -> list:
    # Dados para BD:
    server = function-to-get-the-value
    database = function-to-get-the-value
    username = function-to-get-the-value
    password = function-to-get-the-value
    driver = 'ODBC Driver 17 for SQL Server'

    con_str='Driver={ODBC Driver 17 for SQL Server};Server=tcp:'+server+',1433;Database='+database+';Uid='+username+';Pwd='+password+';Encrypt=yes;TrustServerCertificate=no;Connection Timeout=30;Authentication=ActiveDirectoryPassword'

    sql = """
    SELECT col1, col2
    FROM table_h
    WHERE date_cre BETWEEN DATEADD(day, -1, SYSDATETIME()) AND SYSDATETIME()
    """

    conn = connect(con_str)
    cursor = conn.cursor()
    cursor.execute(sql)
    rows = cursor.fetchall()
    products_orc = DataFrame.from_records(rows, columns=[col[0] for col in cursor.description])
    
    conn.commit()
    cursor.close()
    conn.close()
    
    return products_orc.to_dict('records')

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 16:01:34