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

Durable Functions中用Pickle传递DataFrame报错,求原因与解决方法

使用Durable Functions实现数据分析时的Pickle序列化报错及解决方法

我尝试用Durable Functions实现基于函数的数据分析,因需传递DataFrame等数据,采用Pickle序列化方式进行数据交换,但运行代码时出现报错,且VS Code发生卡顿,代码及报错截图如下:

import azure.functions as func
import azure.durable_functions as df
import pandas as pd
from sklearn.linear_model import LinearRegression
from sklearn.datasets import fetch_california_housing  # Dataset
from sklearn.model_selection import train_test_split 
from sklearn.linear_model import Lasso  
from sklearn.linear_model import Ridge 
from sklearn.metrics import mean_squared_error  # MSE(Mean Squared Error)
from sklearn.preprocessing import StandardScaler 

app = df.DFApp(http_auth_level=func.AuthLevel.ANONYMOUS)
### client function ###
@app.route(route="orchestrators/client_function")
@app.durable_client_input(client_name="client")
async def client_function(req: func.HttpRequest, client: df.DurableOrchestrationClient) -> func.HttpResponse:
    instance_id = await client.start_new("orchestrator", None, {})
    await client.wait_for_completion_or_create_check_status_response(req, instance_id)
    return client.create_check_status_response(req, instance_id)

### orchestrator function ###
@app.orchestration_trigger(context_name="context")
def orchestrator(context: df.DurableOrchestrationContext) -> str:
    data = yield context.call_activity("prepare_data", '')
    simple = yield context.call_activity("simple_regression", {"data": data})
    multiple = yield context.call_activity("multiple_regression", {"data": data})
    return "finished"

### activity function ###
@app.activity_trigger(input_name="blank")
def prepare_data(blank: str):
    # prepare data
    california_housing = fetch_california_housing()
    exp_data = pd.DataFrame(california_housing.data, columns=california_housing.feature_names) # 説明変数
    tar_data = pd.DataFrame(california_housing.target, columns=['HousingPrices']) # 目的変数
    data = pd.concat([exp_data, tar_data], axis=1) # データを結合

    # Delete anomalous values
    data = data[data['HouseAge'] != 52]
    data = data[data['HousingPrices'] != 5.00001]

    # Create useful variables
    data['Household'] = data['Population']/data['AveOccup']
    data['AllRooms'] = data['AveRooms']*data['Household']
    data['AllBedrms'] = data['AveBedrms']*data['Household']

    data = pickle.dumps(data)
    return data

### simple regression analysis ###
@app.activity_trigger(input_name="arg")
def simple_regression(arg: dict):
    data = pickle.loads(arg['data'])
    exp_var = 'MedInc'
    tar_var = 'HousingPrices'

    # Remove outliers
    q_95 = data[exp_var].quantile(0.95)
    data = data[data[exp_var] < q_95]

    # Split data into explanatory and objective variables
    X = data[[exp_var]]
    y = data[[tar_var]]

    # learn
    model = LinearRegression()
    model.fit(X, y)

    model = pickle.dumps(model)
    return model

### multiple regression analysis ###
@app.activity_trigger(input_name="arg")
def multiple_regression(arg: dict):
    data = pickle.loads(arg['data'])
    exp_vars = ['MedInc', 'HouseAge', 'AveRooms', 'AveBedrms', 'Population', 'AveOccup', 'Latitude', 'Longitude']
    tar_var = 'HousingPrices'

    # Remove outliers
    for exp_var in exp_vars:
        q_95 = data[exp_var].quantile(0.95)
        data = data[data[exp_var] < q_95]

    # Split data into explanatory and objective variables
    X = data[exp_vars]
    y = data[[tar_var]]

    # Split into training and test data
    X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.2, random_state=0)

    #  Standardize X_train
    scaler = StandardScaler()
    scaler.fit(X_train)
    X_train_scaled = scaler.transform(X_train)
    X_train_scaled = pd.DataFrame(X_train_scaled, columns = exp_vars)

    # learn
    model = LinearRegression()
    model.fit(X_train_scaled, y_train)

    model = pickle.dumps(model)
    X_train_scaled = pickle.dumps(X_train_scaled)
    y_train = pickle.dumps(y_train)
    X_test = pickle.dumps(X_test)
    y_test = pickle.dumps(y_test)
    scaler = pickle.dumps(scaler)
    return model, X_train_scaled, y_train, X_test, y_test, scaler

报错截图

报错原因

  • 代码未导入pickle模块,直接调用pickle.dumps()/pickle.loads()触发NameError。
  • Durable Functions默认依赖JSON序列化传递活动函数间的数据,Pickle序列化后的字节流(bytes类型)无法被JSON直接序列化,导致异常。
  • multiple_regression函数返回多个Pickle序列化对象组成的元组,加重序列化处理负担;同时客户端函数使用await client.wait_for_completion_or_create_check_status_response阻塞请求直到Orchestrator完成,数据处理耗时较长时会导致VS Code卡顿。

解决方法

1. 导入缺失模块

在代码顶部添加pickle和base64模块:

import pickle
import base64

2. 对Pickle字节流做Base64编码适配JSON

将Pickle生成的bytes转成Base64字符串传递,接收时解码后再反序列化:

  • 修改prepare_data的返回逻辑:
# 替换原有return部分
data_pickle = pickle.dumps(data)
return base64.b64encode(data_pickle).decode('utf-8')
  • 修改simple_regression和multiple_regression的数据加载逻辑:
# 替换原有data加载部分
data_pickle = base64.b64decode(arg['data'])
data = pickle.loads(data_pickle)
  • 对multiple_regression返回的每个序列化对象也做Base64编码,并用字典包装返回值:
# 替换原有return部分
return {
    "model": base64.b64encode(model).decode('utf-8'),
    "X_train_scaled": base64.b64encode(X_train_scaled).decode('utf-8'),
    "y_train": base64.b64encode(y_train).decode('utf-8'),
    "X_test": base64.b64encode(X_test).decode('utf-8'),
    "y_test": base64.b64encode(y_test).decode('utf-8'),
    "scaler": base64.b64encode(scaler).decode('utf-8')
}

3. 优化客户端函数避免卡顿

移除阻塞式等待,直接返回状态检查响应,通过轮询获取结果:

async def client_function(req: func.HttpRequest, client: df.DurableOrchestrationClient) -> func.HttpResponse:
    instance_id = await client.start_new("orchestrator", None, {})
    return client.create_check_status_response(req, instance_id)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 17:04:52