【问题标题】:Why PySpark execute only the default statement in my custom `SQLTransformer`为什么 PySpark 在我的自定义“SQLTransformer”中只执行默认语句
【发布时间】:2019-04-16 07:46:42
【问题描述】:

我在 PySpark 中编写了一个自定义 SQLTransformer。并且必须设置默认 SQL 语句才能执行代码。我可以在 Python 中保存 custum 转换器,加载它并使用 Scala 或/和 Python 执行它,但是尽管_transform 方法中还有其他内容,但只执行默认语句。我对两种语言都有相同的结果,那么问题与_to_java 方法或JavaTransformer 类无关。

class filter(SQLTransformer): 
    def __init__(self):
        super(filter, self).__init__() 
        self._setDefault(statement = "select text, label from __THIS__") 

    def _transform(self, df): 
        df = df.filter(df.id > 23)
        return df

【问题讨论】:

  • 我需要在 Scala 管道中调用 SQLTransformer。我可以在 Python 中保存 SQLTransformer,在 Scala 端加载并运行它,但尽管我在类中定义了 _transform 方法,但默认语句在 Scala 端执行。

标签: apache-spark pyspark pipeline apache-spark-ml


【解决方案1】:

不支持此类信息流。要创建可与 Python 和 Scala 代码库一起使用的 Tranformer,您需要:

  • 实现 Java 或 Scala Transformer,在您的情况下扩展 org.apache.spark.ml.feature.SQLTransformer
  • 以与pyspark.sql.ml.feature.SQLTransformer 相同的方式添加扩展pyspark.sql.ml.wrapper.JavaTransformer 的Python 包装器,并与它对应的JVM 接口。

【讨论】:

  • 谢谢,也就是说用 Python 编写的自定义 Transformer 不能在 Scala 管道中使用。因为,如果我需要在 Scala 和 Python 中编写相同的代码,最好直接在我的 Scala 管道中使用已经用 Scala 编写的代码。
猜你喜欢
  • 1970-01-01
  • 2017-07-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-09-19
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多