如何在sparklyr中按组分布式训练并拟合H2O模型
解决Spark + H2O按组拟合Cox模型的问题
我来帮你搞定这个问题——你遇到的核心卡点是driver端的SparkContext(sc)没办法序列化传递到worker节点。spark_apply的匿名函数是在每个worker的task里执行的,driver端的sc属于上下文对象,根本没法被序列化发送到远端,这就是你看到Spark任务失败的直接原因。
下面是修正后的可行方案,还有几个关键注意事项:
1. 修正spark_apply代码:移除sc参数,让worker自行初始化H2O
不用在spark_apply的函数里传sc,而是让每个worker任务自己处理H2O的初始化,同时正确转换分组数据为H2O Frame:
# 先确保你已经正确配置了Sparkling Water环境 library(sparklyr) library(h2o) library(sparklingwater) # Driver端初始化Spark和H2O整合(只在driver跑一次) sc <- spark_connect(master = "local[*]") h2o_context(sc) # 初始化Sparkling Water上下文 # 修正后的按组建模代码 df %>% spark_apply( function(e) { # Worker节点初始化H2O(如果还没初始化) if (!h2o_initialized()) { # 用随机端口避免多worker端口冲突 h2o.init(ip = "localhost", port = 54321 + sample(1:100, 1), startH2O = TRUE) } # 转换分组数据为H2O Frame,不需要传sc! h2o_df <- as.h2o(e, strict_version_check = FALSE) # 拟合Cox比例风险模型 model <- h2o.coxph( x = predictors, event_column = "event", stop_column = "time_to_next", training_frame = h2o_df ) # 注意:别直接返回H2O模型对象(序列化会出问题),返回你需要的关键结果 data.frame( id = unique(e$id), model_coefficients = list(h2o.coef(model)), aic = h2o.aic(model), stringsAsFactors = FALSE ) }, group_by = "id" )
2. 必须注意的几个关键点
- 绝对不要传递driver端对象到worker:sc、h2o_context这类driver专属的上下文对象,绝对不能在spark_apply的函数里引用,worker必须用自己本地的资源初始化。
- 返回可序列化的结果:H2O模型对象本身很难序列化,所以一定要提取你需要的信息(比如系数、AIC、预测结果等)转成普通DataFrame再返回。
- 避免端口冲突:多个worker同时启动H2O实例时,要随机指定端口,不然会出现端口占用的问题,上面的代码已经帮你处理了这个点。
- 环境版本一致:所有worker节点的R环境必须安装和driver完全相同版本的h2o、sparklyr、sparklingwater包,版本不匹配大概率会导致任务失败。
3. 小数据集备选方案:本地按组建模
如果你的数据量不大,也可以先把数据拉到本地,用H2O结合dplyr的分组功能建模,这种方式更简单,但只适合小数据集:
# 把Spark DataFrame拉到本地 local_df <- df %>% collect() # 按id分组拟合模型 model_list <- local_df %>% group_by(id) %>% group_map(function(group_data, group_key) { h2o_df <- as.h2o(group_data) h2o.coxph( x = predictors, event_column = "event", stop_column = "time_to_next", training_frame = h2o_df ) })
不过这种方式没办法分布式处理,大数据场景还是推荐第一种spark_apply的方案。
内容的提问来源于stack exchange,提问作者DataTx
相关产品推荐
相关产品推荐

