【发布时间】:2020-03-23 22:44:37
【问题描述】:
在转换使用 VectorAssembler 的 ML Pipeline 时,会遇到“Param handleInvalid does not exist”错误。为什么会这样?我错过了什么吗?我是 PySpark 的新手。
我根据代码使用它来将给定的列列表组合成一个向量列:
for categoricalCol in categoricalColumns:
stringIndexer = StringIndexer(inputCol = categoricalCol, outputCol = categoricalCol + 'Index').setHandleInvalid("keep")
encoder = OneHotEncoderEstimator(inputCols=[stringIndexer.getOutputCol()], outputCols=[categoricalCol + "classVec"])
stages += [stringIndexer, encoder]
label_stringIdx = StringIndexer(inputCol = 'response', outputCol = 'label')
stages += [label_stringIdx]
numericCols = ['date_acct_', 'date_loan_', 'amount', 'duration', 'payments', 'birth_number_', 'min1', 'max1', 'mean1', 'min2', 'max2', 'mean2', 'min3', 'max3', 'mean3', 'min4', 'max4', 'mean4', 'min5', 'max5', 'mean5', 'min6', 'max6', 'mean6', 'gen', 'has_card']
assemblerInputs = [c + "classVec" for c in categoricalColumns] + numericCols
assembler = VectorAssembler(inputCols=assemblerInputs, outputCol="feature")
print(assembler)
stages += [assembler]
df_features 是我保存所有列的主要数据框。我试图保留 handleInvalid = 'keep' 和 handleInvalid = 'skip' 但不幸的是得到了同样的错误。
得到以下错误:
Traceback (most recent call last):
File "spark_model_exp_.py", line 275, in <module>
feature_df = assembler.transform(features)
File "/usr/local/lib/python3.6/site-packages/pyspark/ml/base.py", line 173, in transform
return self._transform(dataset)
File "/usr/local/lib/python3.6/site-packages/pyspark/ml/wrapper.py", line 311, in _transform
self._transfer_params_to_java()
File "/usr/local/lib/python3.6/site-packages/pyspark/ml/wrapper.py", line 124, in _transfer_params_to_java
pair = self._make_java_param_pair(param, self._paramMap[param])
File "/usr/local/lib/python3.6/site-packages/pyspark/ml/wrapper.py", line 113, in _make_java_param_pair
java_param = self._java_obj.getParam(param.name)
File "/usr/local/lib/python3.6/site-packages/py4j/java_gateway.py", line 1257, in __call__
answer, self.gateway_client, self.target_id, self.name)
File "/usr/local/lib/python3.6/site-packages/pyspark/sql/utils.py", line 63, in deco
return f(*a, **kw)
File "/usr/local/lib/python3.6/site-packages/py4j/protocol.py", line 328, in get_return_value
format(target_id, ".", name), value)
py4j.protocol.Py4JJavaError: An error occurred while calling o1072.getParam.
: java.util.NoSuchElementException: Param handleInvalid does not exist.
at org.apache.spark.ml.param.Params$$anonfun$getParam$2.apply(params.scala:729)
at org.apache.spark.ml.param.Params$$anonfun$getParam$2.apply(params.scala:729)
at scala.Option.getOrElse(Option.scala:121)
at org.apache.spark.ml.param.Params$class.getParam(params.scala:728)
at org.apache.spark.ml.PipelineStage.getParam(Pipeline.scala:43)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.lang.reflect.Method.invoke(Method.java:498)
at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
at py4j.Gateway.invoke(Gateway.java:282)
at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
at py4j.commands.CallCommand.execute(CallCommand.java:79)
at py4j.GatewayConnection.run(GatewayConnection.java:238)
at java.lang.Thread.run(Thread.java:745)
我之前尝试过什么?
categoricalColumns = ['frequency', 'type_disp', 'type_card']
for categoricalCol in categoricalColumns:
stringIndexer = StringIndexer(inputCol = categoricalCol, outputCol = categoricalCol + 'Index').setHandleInvalid("keep")
encoder = OneHotEncoderEstimator(inputCols=[stringIndexer.getOutputCol()], outputCols=[categoricalCol + "classVec"])
stages += [stringIndexer, encoder]
label_stringIdx = StringIndexer(inputCol = 'response', outputCol = 'label')
stages += [label_stringIdx]
numericCols = ['date_acct_', 'date_loan_', 'amount', 'duration', 'payments', 'birth_number_', 'min1', 'max1', 'mean1', 'min2', 'max2', 'mean2', 'min3', 'max3', 'mean3', 'min4', 'max4', 'mean4', 'gen', 'has_card']
assemblerInputs = [c + "classVec" for c in categoricalColumns] + numericCols
assembler = VectorAssembler(inputCols=assemblerInputs, outputCol="feature")
stages += [assembler]
pipeline = Pipeline(stages = stages)
pipelineModel = pipeline.fit(features)
features = pipelineModel.transform(features)
features.show(n=2)
selectedCols = ['label', 'feature'] + cols
features = features.select(selectedCols)
print(features.dtypes)
在上面的代码含义中,也使用了 Pipeline,我在 Pipeline 的转换功能处遇到了错误。当我尝试上面的代码时,我在 VectorAssembler 转换函数中没有收到错误,在 Pipeline 转换函数中也没有收到相同的错误(Param handleInvalid 不存在)。
请让我知道这方面的更多细节。我们可以尝试通过其他一些替代方案来实现这一目标吗?
编辑:我得到了为什么会发生这种情况的部分答案,因为在本地 spark 版本 = 2.4 上,所以代码在这个上运行良好,但集群 spark 版本 = 2.3 并且由于 handleInvalid 是从版本 2.4 引入的因此我收到此错误。
但我想知道,因为我已经检查过数据帧中没有 NULL/NaN 值,但是 vectorAssembler 是如何调用 handleInvalid 参数的?我在想我是否可以绕过这个隐式调用 handleInvalid 以便我不应该面对这个错误,或者是否有任何其他替代选项而不是将 spark 版本从 2.3 升级到 2.4?
有人可以就此提出建议吗?
【问题讨论】:
-
能否提供您的 Spark 版本?
-
@napoleon_borntoparty Spark 版本-2.3.2.3.1.4.0-315 请让我知道更多详情。
-
我想在这里提供的另一个输入是,当我在 edgenode 上本地使用 spark 时,我能够生成完整的模型,但是当我尝试在集群的基础上运行我的模型时,它就坏了解决我提到的错误。
标签: apache-spark pyspark apache-spark-mllib apache-spark-ml apache-spark-2.0