Azure Durable 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

