【发布时间】: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