Azure Durable Function(Python)报Worker加载函数失败错误
Azure Durable Function 函数加载报错根因与修复方案
根因分析
报错是三个核心用法不符合Azure Durable Function开发规范导致的:
- Azure Functions Python工作器加载函数时,会强制校验
main入口的所有入参,必须和对应function.json中声明的绑定项一一对应。你给DMBucket1等activity函数的main方法直接定义了df、df_curr等7个普通入参,这些参数既不是官方要求的触发器绑定对象,也没有在function.json里声明对应输入绑定,工作器在函数加载阶段就会直接拦截抛出FunctionLoadError,函数根本没机会启动执行。 - Durable Function的
context.call_activity方法仅支持传入1个可序列化的载荷参数,你当前代码里给call_activity传了多个位置参数,后续参数会被运行时直接丢弃,就算过了加载校验,调用逻辑也会出错。 - 额外隐患:pandas DataFrame对象无法直接被Durable Task框架默认的JSON序列化器正确处理,直接跨activity传递会出现类型丢失、序列化失败的问题。
修复方案
按以下步骤逐一调整代码即可解决问题:
- 统一所有activity函数的入口签名
所有activity函数的main方法只保留1个入参,用来接收orchestrator传过来的载荷,禁止自定义多个普通入参。
以DMPullFromDB为例,调整后的代码如下:import logging import pandas as pd def preprocessing(df): df=df.applymap(lambda x: x.lower() if type(x) == str else x) return df def name_fix(df): df1=df.groupby(["employeeid"]).first().reset_index() return(df1) # 只保留1个入参接收传入载荷 def main(input_param): df_roster=pd.read_csv("df_roster.csv") df_curr=pd.read_csv("df_curr.csv") df_prev=pd.read_csv("df_prev.csv") df=pd.read_csv("df.csv") manager_name=pd.read_csv("manager_name.csv") name_data=pd.read_csv("name_data.csv") df_roster=preprocessing(df_roster) df_curr=preprocessing(df_curr) df_prev=preprocessing(df_prev) df=preprocessing(df) manager_name=preprocessing(manager_name) name_data=preprocessing(name_data) logging.info('Preprocessing is complete') name_data1=name_fix(name_data) df_roster['recorddate'] = pd.to_datetime(df_roster['recorddate']) df_roster['recmonth'] = df_roster['recorddate'].dt.month df_roster['recyear'] = df_roster['recorddate'].dt.year data=df.merge(df_roster,on = ['companyguid','organization','primaryprogram','employeeid','recyear','recmonth'],how="left") data1 = data[['companyguid','organization','primaryprogram','employeeid','recmonth','recyear','metricname','metrictype','goal','actualvalue','Manageremployeeid']] # 返回时把DataFrame转成可JSON序列化的字典格式,不要直接返回DataFrame对象 return { "df": df.to_dict(orient="records"), "df_curr": df_curr.to_dict(orient="records"), "df_prev": df_prev.to_dict(orient="records"), "df_roster": df_roster.to_dict(orient="records"), "name_data1": name_data1.to_dict(orient="records"), "manager_name": manager_name.to_dict(orient="records"), "data1": data1.to_dict(orient="records") } - 调整下游activity函数的参数解析逻辑
DMBucket1、DMBucket2、DMBucket3、DMBucket4、DMRandomization、DMPushToDB所有activity的main方法都改成单入参,拿到传入的字典载荷后,再自行把字典转回pandas DataFrame执行业务逻辑。
以DMBucket1为例:import pandas as pd def main(input_payload): # 把传入的字典转回DataFrame df = pd.DataFrame(input_payload["df"]) df_curr = pd.DataFrame(input_payload["df_curr"]) df_prev = pd.DataFrame(input_payload["df_prev"]) df_roster = pd.DataFrame(input_payload["df_roster"]) name_data1 = pd.DataFrame(input_payload["name_data1"]) manager_name = pd.DataFrame(input_payload["manager_name"]) data1 = pd.DataFrame(input_payload["data1"]) # 原有DMBucket1的业务逻辑写在这里 # 最终返回结果同样转成可JSON序列化的格式(字典/数值/字符串均可,不要直接返回DataFrame) return bucket1_result - 修正orchestrator的activity调用逻辑
调用call_activity时,把需要传给下游的所有参数打包成1个字典作为载荷传入,禁止传多个位置参数,调整后的orchestrator代码如下:import logging import azure.durable_functions as df def orchestrator_function(context: df.DurableOrchestrationContext): db_pull_result = yield context.call_activity('DMPullFromDB', None) # 打包传给bucket类函数的参数 bucket_input = { "df": db_pull_result["df"], "df_curr": db_pull_result["df_curr"], "df_prev": db_pull_result["df_prev"], "df_roster": db_pull_result["df_roster"], "name_data1": db_pull_result["name_data1"], "manager_name": db_pull_result["manager_name"], "data1": db_pull_result["data1"] } bucket1 = yield context.call_activity('DMBucket1', bucket_input) bucket2 = yield context.call_activity('DMBucket2', bucket_input) bucket3 = yield context.call_activity('DMBucket3', bucket_input) bucket4 = yield context.call_activity('DMBucket4', bucket_input) # 打包传给随机化函数的参数 rand_input = { "bucket1": bucket1, "bucket2": bucket2, "bucket3": bucket3, "bucket4": bucket4, "data1": db_pull_result["data1"], "df_roster": db_pull_result["df_roster"] } rand = yield context.call_activity('DMRandomization', rand_input) pushDB = yield context.call_activity('DMPushToDB', rand) if pushDB==1: return("The DB push was successful") return("DB push was not successful") main = df.Orchestrator.create(orchestrator_function) - 校验function.json配置
检查每个activity对应的function.json,确保bindings数组中仅保留activity trigger配置,不需要额外声明其他输入绑定。
注意:如果处理的DataFrame数据量较大,不建议直接转字典传递,会导致序列化后的消息体积超过Azure Storage Queue单消息大小上限,这种场景可以先把DataFrame序列化后存入临时blob,仅把blob路径作为载荷传给下游activity,由下游自行读取blob获取数据即可。
内容的提问来源于stack exchange,提问作者AGS
相关产品推荐
相关产品推荐

