【问题标题】:How to load a custom transformer in Spark 2.4如何在 Spark 2.4 中加载自定义转换器
【发布时间】:2019-04-21 02:35:58
【问题描述】:

我正在尝试在 Spark 2.4.0 中创建自定义转换器。保存它工作正常。但是,当我尝试加载它时,我收到以下错误:

java.lang.NoSuchMethodException: TestTransformer.<init>(java.lang.String)
  at java.lang.Class.getConstructor0(Class.java:3082)
  at java.lang.Class.getConstructor(Class.java:1825)
  at org.apache.spark.ml.util.DefaultParamsReader.load(ReadWrite.scala:496)
  at org.apache.spark.ml.util.MLReadable$class.load(ReadWrite.scala:380)
  at TestTransformer$.load(<console>:40)
  ... 31 elided

这表明它找不到我的转换器的构造函数,这对我来说真的没有意义。

MCVE:

import org.apache.spark.sql.{Dataset, DataFrame}
import org.apache.spark.sql.types.{StructType}
import org.apache.spark.ml.Transformer
import org.apache.spark.ml.param.ParamMap
import org.apache.spark.ml.util.{DefaultParamsReadable, DefaultParamsWritable, Identifiable}

class TestTransformer(override val uid: String) extends Transformer with DefaultParamsWritable{

    def this() = this(Identifiable.randomUID("TestTransformer"))

    override def transform(df: Dataset[_]): DataFrame = {
        val columns = df.columns
        df.select(columns.head, columns.tail: _*)
    }

    override def transformSchema(schema: StructType): StructType = {
        schema
    }

    override def copy(extra: ParamMap): TestTransformer = defaultCopy[TestTransformer](extra)
}

object TestTransformer extends DefaultParamsReadable[TestTransformer]{

    override def load(path: String): TestTransformer = super.load(path)

}

val transformer = new TestTransformer("test")

transformer.write.overwrite().save("test_transformer")
TestTransformer.load("test_transformer")

运行此程序(我使用的是 Jupyter 笔记本)会导致上述错误。我尝试将其编译为 .jar 文件并将其运行,没有任何区别。

让我感到困惑的是,等效的 PySpark 代码运行良好:

from pyspark.sql import SparkSession, DataFrame
from pyspark.ml import Transformer
from pyspark.ml.util import DefaultParamsReadable, DefaultParamsWritable

class TestTransformer(Transformer, DefaultParamsWritable, DefaultParamsReadable):

    def transform(self, df: DataFrame) -> DataFrame:
        return df

TestTransformer().save('test_transformer')
TestTransformer.load('test_transformer')

如何制作可以保存和加载的自定义 Spark 转换器?

【问题讨论】:

    标签: java scala apache-spark


    【解决方案1】:

    我可以在 spark-shell 中重现您的问题。

    为了找到问题的根源,我查看了 DefaultParamsReadableDefaultParamsReader 源,我可以看到它们利用了 Java 反射。

    https://github.com/apache/spark/blob/v2.4.0/mllib/src/main/scala/org/apache/spark/ml/util/ReadWrite.scala

    第 495-496 行

    val instance =
        cls.getConstructor(classOf[String]).newInstance(metadata.uid).asInstanceOf[Params]
    

    我认为 scala REPL 和 Java 反射不是好朋友。

    如果你运行这个 sn-p(在你的之后):

    new TestTransformer().getClass.getConstructors
    

    你会得到以下输出:

    res1: Array[java.lang.reflect.Constructor[_]] = Array(public TestTransformer($iw), public TestTransformer($iw,java.lang.String))
    

    这是真的! TestTransformer.&lt;init&gt;(java.lang.String) 不存在。

    我找到了 2 个解决方法,

    1. 用 sbt 编译你的代码并创建一个 jar,然后用 :require 包含在 spark-shell 中,对我有用(你提到你尝试了一个 jar,但我不知道如何)

    2. 使用 :paste -raw 将代码粘贴到 spark-shell 中,效果也很好。我想-raw 会阻止 REPL 对你的班级进行恶作剧。 见:https://docs.scala-lang.org/overviews/repl/overview.html

    我不确定如何将其中任何一个调整到 Jupyter,但我希望这些信息对你有用。

    注意:我实际上在 spark 2.4.1 中使用了 spark-shell

    【讨论】:

    • 我对 Scala 还是很陌生,所以也许我在 .jar 中并没有真正做到这一点(我对 Python 和 PySpark 有更多的经验)。关于你的回答,我明天试试,谢谢!我不确定这是否超出了这个问题的范围,但你知道这是否也适用于 Spark 2.2 和 2.3?
    • 我所做的是创建一个最小的 sbt 项目并使用package 命令创建了 jar。是的,我非常有信心 2.2 和 2.3 的问题和解决方案是相同的。事实上,我在 spark 2.3.2 中测试了上述部分
    猜你喜欢
    • 1970-01-01
    • 2020-08-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-09-28
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多