【问题标题】:Caching intermediate results in Spark ML pipeline在 Spark ML 管道中缓存中间结果
【发布时间】:2016-09-02 15:47:58
【问题描述】:

最近我计划将我的独立 python ML 代码迁移到 spark。 spark.ml 中的 ML 管道非常方便,具有用于链接算法阶段和超参数网格搜索的简化 API。

不过,我发现它对现有文档中的一项重要功能的支持并不明显:缓存中间结果。当流水线涉及计算密集阶段时,此功能的重要性就显现出来了。

例如,在我的例子中,我使用一个巨大的稀疏矩阵对时间序列数据执行多个移动平均值,以形成输入特征。矩阵的结构由一些超参数决定。这一步结果成为整个管道的瓶颈,因为我必须在运行时构造矩阵。

在参数搜索过程中,除了这个“结构参数”之外,我通常还有其他参数要检查。因此,如果我可以在“结构参数”不变的情况下重用巨大的矩阵,我可以节省大量时间。出于这个原因,我特意编写了代码来缓存和重用这些中间结果。

所以我的问题是:Spark 的 ML 管道能否自动处理中间缓存?还是我必须手动形成代码才能这样做?如果是这样,是否有任何最佳实践可供学习?

附:我查看了官方文档和其他一些材料,但似乎都没有讨论这个话题。

【问题讨论】:

    标签: apache-spark apache-spark-ml


    【解决方案1】:

    所以我遇到了同样的问题,我解决的方法是我实现了自己的 PipelineStage,它缓存输入 DataSet 并按原样返回。

    import org.apache.spark.ml.Transformer
    import org.apache.spark.ml.param.ParamMap
    import org.apache.spark.ml.util.{DefaultParamsWritable, Identifiable}
    import org.apache.spark.sql.{DataFrame, Dataset}
    import org.apache.spark.sql.types.StructType
    
    class Cacher(val uid: String) extends Transformer with DefaultParamsWritable {
      override def transform(dataset: Dataset[_]): DataFrame = dataset.toDF.cache()
    
      override def copy(extra: ParamMap): Transformer = defaultCopy(extra)
    
      override def transformSchema(schema: StructType): StructType = schema
    
      def this() = this(Identifiable.randomUID("CacherTransformer"))
    }
    

    要使用它,你可以这样做:

    new Pipeline().setStages(Array(stage1, new Cacher(), stage2))
    

    【讨论】:

    • 我想唯一的问题是解决方案(我实际上赞成!)是您不会取消保留先前缓存的数据帧(如果您链接多个缓存器)。你可能会争辩说这不是问题,因为 Spark 在 GC 时会自动取消持久化,但它可能会让你的 UI 变得非常混乱,例如看到这么多缓存数据。
    猜你喜欢
    • 1970-01-01
    • 2020-02-16
    • 2016-05-23
    • 1970-01-01
    • 2017-08-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-03-29
    相关资源
    最近更新 更多