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

增加数据量后PySpark Pandas UDF返回错误的问题求助

增加数据量后PySpark Pandas UDF返回错误的问题求助

各位大佬好,我刚接触PySpark,最近写了一个Pandas UDF用来接收时间序列数据并训练不同的机器学习模型,但遇到了头疼的问题——数据量小的时候一切正常,可数据量增大后,这个UDF就开始返回错误了。先跟大家梳理下我的处理流程和代码,麻烦帮忙看看哪里出问题了:

我的原始数据包含这些列:ID(产品ID,字符串类型)、models_name(模型名称,比如LR、XGboost、SARIMAX等)、prices、units_sold(这两列是时间序列数据)还有date列。

为了把时间序列转换成列表格式传给UDF,我先做了分组聚合:

grouped_spark_df = spark_df.groupBy("ID").agg(
                    F.collect_list("prices").alias("prices"),
                    F.collect_list("units_sold").alias("units_sold"),
                    F.collect_list("date").alias("date"))

分组后的数据集长这样:

IDpricesunits_solddate
P_1[10.5, 11.0, 9.8][2, 3, 1]["2024-03-01", "2024-03-02", "2024-03-03"]
P_2[15.0, 16.2][5, 7]["2024-03-01", "2024-03-02"]
P_3[7.5][10]["2024-03-01"]

接下来我需要给每个产品匹配所有要训练的模型,所以用explode把模型列表展开成列:

exploded_df = grouped_spark_df.withColumn("models_name", explode(array([lit(m) for m in models_list])))

处理后的数据集变成了这样:

IDpricesunits_solddatemodels_name
P_1[10.5, 11.0, 9.8][2, 3, 1]["2024-03-01", "2024-03-02", "2024-03-03"]A
P_1[10.5, 11.0, 9.8][2, 3, 1]["2024-03-01", "2024-03-02", "2024-03-03"]B
P_1[10.5, 11.0, 9.8][2, 3, 1]["2024-03-01", "2024-03-02", "2024-03-03"]C
P_2[15.0, 16.2][5, 7]["2024-03-01", "2024-03-02"]A
P_2[15.0, 16.2][5, 7]["2024-03-01", "2024-03-02"]B
P_2[15.0, 16.2][5, 7]["2024-03-01", "2024-03-02"]C
P_3[7.5][10]["2024-03-01"]A
P_3[7.5][10]["2024-03-01"]B
P_3[7.5][10]["2024-03-01"]C

之后我对数据做了分区,然后调用自定义的Pandas UDF:

partitioned_df = exploded_df.repartition(50, "Bucket","ID","models_name").coalesce(10)
results_spark = partitioned_df.select(
        col("ID"),
        col("models_name"),
        train_model_udf(
            col("ID"), col("prices"), col("units_sold"), col("date")
         
        ).alias("model_metrics")
    )

最后是我的Pandas UDF代码,里面包含了数据转换、模型训练、结果保存等逻辑:

from pyspark.sql.functions import pandas_udf
from pyspark.sql.types import MapType, StringType, FloatType,ArrayType
import pandas as pd
from sklearn.model_selection import train_test_split

 
def create_broadcast_and_udf(feat_cols,model_results_data,feature_set_name,target,dt_string):
    # broadcast variables....
    @pandas_udf(StringType())
   
    def train_model_udf(
            ID, prices, units_sold, date, models_name):

        # Convert to Pandas DataFrame
            data = pd.DataFrame({
                "ID": ID, "prices": prices, "units_sold": units_sold, "date": date, 
                               "models_name": models_name
            })
            for col in data.columns:
                data[col] = data[col].apply(lambda x: x.tolist() if isinstance(x, np.ndarray) else x)

            explode_cols = [col for col in data.columns if col != "ID" and isinstance(data[col][0], list)]

        # Explode all selected columns
            data = data.explode(explode_cols, ignore_index=True)
     
            try:
                
                new_data = []
                ID = str(data["ID"].iloc[0])
                # return  pd.Series([models_name] * len(ID))
                model_directory = f'{model_results_data}/{dt_string}/{models_name}/{feature_set_name}_{target}'
                if not os.path.exists(model_directory): # checking if model directory exists or not
                                os.makedirs(model_directory)
                
   
                train, test = train_test_split(data, test_size=0.2, shuffle=False) 
                combined_data = pd.concat([train, test]).sort_values(by=['ID', 'WEEK'])

                data.to_csv(f"{model_directory}/{ID}_ID_CombinedData.csv", index=False)
                model_write_path = f"{model_directory}/{ID}.pkl" 
                                        
                # training models
                if models_name.lower() =='sarimax':
                    elasticity, y_train_pred, forecast ,mape_hj,wmape_hj,rmse,mse,mae  = train_sarimax(train,test,model_write_path,price_feature,feat_cols,target,elasticity_method)
                elif models_name.lower() =='lr':
                    elasticity, y_train_pred, forecast ,mape_hj,wmape_hj,rmse,mse,mae  = train_lr(train,test,model_write_path,price_feature,broadcast_features.value,target)
                elif models_name.lower() =='xgboost':
                    elasticity, y_train_pred, forecast ,mape_hj,wmape_hj,rmse,mse,mae   = train_xgb(train,test,model_write_path,price_feature,feat_cols,target)
                elif models_name.lower() =='gam':
                    spline_features = feat_cols.copy()
                    elasticity, y_train_pred, forecast ,mape_hj,wmape_hj,rmse,mse,mae  = train_gam(data,test,ID,spline_features,price_feature,target,model_write_path)
                    print(price_feature)
                    print(target)
                temp_df = data.copy()
                predicted_volumes =   np.concatenate((y_train_pred, forecast)) 
                # Add new columns to the copy
                temp_df['Experiment_Model'] = models_name
                temp_df['feat_cols'] =  "-".join(feat_cols)
                temp_df['MAPE'] = mape_hj
                temp_df['WMAPE'] = wmape_hj
                temp_df['final_elasticity'] = elasticity
                temp_df['predicted_volume'] = predicted_volumes
                temp_df = temp_df.reset_index()
                # Append the rows of the modified DataFrame to the list
                new_data.extend(temp_df.to_dict('records'))

                # calculation percent price for calculating simulated volume 
                test['pct_change_pr'] = test['updated_WAP'].pct_change()

                # Initializing list to store simulated volumes
                simulated_volumes = []

                # Start with the first available volume in the test set as the initial simulated volume
                initial_volume = test[target].iloc[0]
                simulated_volumes.append(initial_volume)

                # Iterate over the test set, starting from the second row
                for i in range(1, len(test)):
                    previous_simulated_volume = simulated_volumes[-1]  # Get the last simulated volume
                    pct_change_pr = test['pct_change_pr'].iloc[i]  # Current row's price change

                    # Calculate the new simulated volume based on the previous simulated volume
                    simulated_volume = previous_simulated_volume + previous_simulated_volume * (elasticity * pct_change_pr)
                    
                    # Append the calculated volume to the list
                    simulated_volumes.append(simulated_volume)

                # Add the simulated volumes to the test set as a new column
                test['simulated_volume'] = simulated_volumes
                test['elasticity'] = elasticity
                                    
                # saving file to check if simulated volume is correctly calculated
                test[['ID','updated_WAP','pct_change_pr',target,'elasticity','simulated_volume']].to_csv(f"{model_directory}/{ID}_simulated.csv")

                # for plotting the graphs 
                temp_test_df = pd.DataFrame()

                temp_test_df['date']= test[target].index

                temp_test_df['actual']= test[target].values
                temp_test_df['predicted']= forecast
                temp_test_df['ppg_name']= ID
                temp_test_df['features used']= [feat_cols]* len(temp_test_df)
                temp_test_df['rmse']= rmse
                temp_test_df['mse']= mse
                # temp_test_df['mae']= mae
                temp_test_df['mape']= mape_hj
                temp_test_df['wmape']= wmape_hj
                temp_test_df['elasticity']= elasticity
                temp_test_df['point_wise_error'] = [abs(g - p) / g * 100 if g != 0 else 0 for g, p in zip(test[target], forecast)]
                test = test.reset_index(drop=True)

                temp_test_df['simulated_volume'] = test['simulated_volume']

                train_temp_df = pd.DataFrame()
                train_temp_df['date'] = train.index
                train_temp_df['actual'] = train[target].values
                train_temp_df['predicted'] = y_train_pred
                train_temp_df['point_wise_error'] = [abs(g - p) / g * 100 if g != 0 else 0 for g, p in zip(train[target].values, y_train_pred)]

备注:内容来源于stack exchange,提问作者Fatima Arshad

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 20:15:26