【问题标题】:Spark 2.0.2 PySpark failing to Import collect_listSpark 2.0.2 PySpark 无法导入 collect_list
【发布时间】:2017-05-05 01:06:30
【问题描述】:

我有一个DataFrame 的形式:

+--------------+------------+----+
|             s|variant_hash|call|
+--------------+------------+----+
|C1046::HG02024|    83779208|   0|
|C1046::HG02025|    83779208|   1|
|C1046::HG02026|    83779208|   0|
|C1047::HG00731|    83779208|   0|
|C1047::HG00732|    83779208|   1
              ...

我希望利用collect_list() 将其转换为:

+--------------------+-------------------------------------+
|                   s|                       feature_vector|
+--------------------+-------------------------------------+
|      C1046::HG02024|[(83779208,   0), (68471259,   2)...]|
+--------------------+-------------------------------------+

其中特征向量列是(variant_hash, call) 形式的元组列表。我正计划利用groupBy 和agg(collect_list()) 来完成此结果,但收到以下错误:

Traceback (most recent call last):
  File "/tmp/ba6a891c-529b-4c75-a76f-8ab20f4377ba/ml_on_vds.py", line 43, in <module>
    vector_df = svc_df.groupBy('s').agg(func.collect_list(('variant_hash', 'call')))
  File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/functions.py", line 39, in _
  File "/usr/lib/spark/python/lib/py4j-0.10.3-src.zip/py4j/java_gateway.py", line 1133, in __call__
  File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 63, in deco
  File "/usr/lib/spark/python/lib/py4j-0.10.3-src.zip/py4j/protocol.py", line 323, in get_return_value
py4j.protocol.Py4JError: An error occurred while calling z:org.apache.spark.sql.functions.collect_list. Trace:
py4j.Py4JException: Method collect_list([class java.util.ArrayList]) does not exist
        at py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:318)
        at py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:339)
        at py4j.Gateway.invoke(Gateway.java:274)
        at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
        at py4j.commands.CallCommand.execute(CallCommand.java:79)
        at py4j.GatewayConnection.run(GatewayConnection.java:214)
        at java.lang.Thread.run(Thread.java:745)

下面的代码显示了我的导入。我认为没有必要在 2.0.2 中导入 HiveContext 和 enableHiveSupport,但我希望这样做可以解决问题。可悲的是,没有运气。有人有解决此导入问题的建议吗?

from pyspark.sql import SparkSession
from pyspark import SparkConf, SparkContext, HiveContext
from pyspark.sql.functions import udf, hash, collect_list
from pyspark.sql.types import *
from hail import *
# Initialize the SparkSession
spark = (SparkSession.builder.appName("PopulationGenomics")
        .config("spark.sql.files.openCostInBytes", "1099511627776")
        .config("spark.sql.files.maxPartitionBytes", "1099511627776")
        .config("spark.hadoop.io.compression.codecs", "org.apache.hadoop.io.compress.DefaultCodec,is.hail.io.compress.BGzipCodec,org.apache.hadoop.io.compress.GzipCodec")
        .enableHiveSupport()
        .getOrCreate())

我正在尝试在 gcloud dataproc 集群上运行此代码。

【问题讨论】:

    标签: apache-spark pyspark google-cloud-dataproc


    【解决方案1】:

    所以它在这一行抛出错误 -

    vector_df = svc_df.groupBy('s').agg(func.collect_list(('variant_hash', 'call')))
    

    您将collect_list 称为func.collect_list 但是您将函数导入为 -

    from pyspark.sql.functions import udf, hash, collect_list

    可能你打算将函数导入为 'func' 之类的

    from pyspark.sql import functions as func,

    【讨论】:

    • 感谢Pushkr 的回复。继续并进行了更正,但仍然收到相同的错误消息。查看我的 IDE(未连接到我的 Spark 集群,但包含应该匹配)时,collect_list 似乎没有包含在 pyspark.sql.functions 中...
    • 你能在 spark CLI 中做一个简单的测试来检查你是否能够导入任何 sql 函数并运行它们吗?
    • 能够从函数库中导入并运行hash 和udf。
    • @mongolol - 你找到解决方案了吗?你能导入 collect_list() 吗?
    • @outlier123 遗憾的是,我们已经从 2.0.2 继续前进,并且不再遇到此问题。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-07-16
    • 1970-01-01
    • 2016-10-01
    • 2023-02-15
    • 2018-11-26
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多