PySpark中为何应避免使用for循环?循环训练模型性能解析
关于Spark中循环训练模型的性能与并行性疑问
背景
我正尝试加速部分数据处理pipeline,但有几个问题未得到确切答案:
- 是否部分for循环根据实现方式是可行的?
- 何时使用循环不会造成显著性能损耗?
我已查阅以下资料:
- David Mudrauskas的相关文章
- Stack Overflow的相关解答
- Spark RDD文档,其中提到:
通常不应使用闭包(如循环或本地定义方法)修改全局状态,Spark不保证此类代码在分布式模式下的行为,仅在本地模式可能偶然运行正常。
具体场景与疑问
我需要针对多个标签列训练一系列Logistic Regression模型,目前有两种实现方式:
方式1:使用for循环遍历训练
from pyspark.ml import Pipeline from pyspark.ml.feature import VectorAssembler from pyspark.ml.classification import LogisticRegression dv = ['y1','y2','y3', ...] models = {} for v in dv: assembler = VectorAssembler(inputCols=feature_cols, outputCol='features') model = LogisticRegression(featuresCol='features',labelCol=v,predictionCol=f'prediction_{v}') pipeline = Pipeline(stages=[assembler,model]) pipe = pipeline.fit(train) models[v] = pipe
方式2:逐个显式训练模型
# y1 assembler = VectorAssembler(inputCols=feature_cols, outputCol='features') model = LogisticRegression(featuresCol='features',labelCol='y1',predictionCol=f'prediction_y1') pipeline = Pipeline(stages=[assembler,model]) pipe = pipeline.fit(train) models['y1'] = pipe # y2 assembler = VectorAssembler(inputCols=feature_cols, outputCol='features') model = LogisticRegression(featuresCol='features',labelCol='y2',predictionCol=f'prediction_y2') pipeline = Pipeline(stages=[assembler,model]) pipe = pipeline.fit(train) models['y2'] = pipe # y3 assembler = VectorAssembler(inputCols=feature_cols, outputCol='features') model = LogisticRegression(featuresCol='features',labelCol='y3',predictionCol=f'prediction_y3') pipeline = Pipeline(stages=[assembler,model]) pipe = pipeline.fit(train) models['y3'] = pipe ...
我的疑问:
- 两种实现方式是否等价?使用for循环是否会比逐个显式训练明显更慢?
- 我了解SparkML有内置并行机制,但担心循环会降低并行性,是否存在更优的并行训练模型方法?当前训练速度很慢,可能操作有误,希望得到解答。
内容的提问来源于stack exchange,提问作者quibbles
相关产品推荐
相关产品推荐

