增加数据量后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"))
分组后的数据集长这样:
| ID | prices | units_sold | date |
|---|---|---|---|
| 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])))
处理后的数据集变成了这样:
| ID | prices | units_sold | date | models_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
相关产品推荐
相关产品推荐

