【问题标题】:PySpark error: AttributeError: 'NoneType' object has no attribute '_jvm'PySpark 错误:AttributeError:“NoneType”对象没有属性“_jvm”
【发布时间】:2017-03-10 21:26:03
【问题描述】:

我有时间戳数据集,格式为

我在 pyspark 中编写了一个 udf 来处理这个数据集并作为键值映射返回。但我收到以下错误消息。

数据集:df_ts_list

+--------------------+
|             ts_list|
+--------------------+
|[1477411200, 1477...|
|[1477238400, 1477...|
|[1477022400, 1477...|
|[1477224000, 1477...|
|[1477256400, 1477...|
|[1477346400, 1476...|
|[1476986400, 1477...|
|[1477321200, 1477...|
|[1477306800, 1477...|
|[1477062000, 1477...|
|[1477249200, 1477...|
|[1477040400, 1477...|
|[1477090800, 1477...|
+--------------------+

Pyspark UDF:

>>> def on_time(ts_list):
...     import sys
...     import os
...     sys.path.append('/usr/lib/python2.7/dist-packages')
...     os.system("sudo apt-get install python-numpy -y")
...     import numpy as np
...     import datetime
...     import time
...     from datetime import timedelta
...     ts = np.array(ts_list)
...     if ts.size == 0:
...             count = 0
...             duration = 0
...             st = time.mktime(datetime.now())
...             ymd = str(datetime.fromtimestamp(st).date())
...     else:
...             ts.sort()
...             one_tag = []
...             start = float(ts[0])
...             for i in range(len(ts)):
...                     if i == (len(ts)) - 1:
...                             end = float(ts[i])
...                             a_round = [start, end]
...                             one_tag.append(a_round)
...                     else:
...                             diff = (datetime.datetime.fromtimestamp(float(ts[i+1])) - datetime.datetime.fromtimestamp(float(ts[i])))
...                             if abs(diff.total_seconds()) > 3600:
...                                     end = float(ts[i])
...                                     a_round = [start, end]
...                                     one_tag.append(a_round)
...                                     start = float(ts[i+1])
...             one_tag = [u for u in one_tag if u[1] - u[0] > 300]
...             count = int(len(one_tag))
...             duration = int(np.diff(one_tag).sum())
...             ymd = str(datetime.datetime.fromtimestamp(time.time()).date())
...     return {'count':count,'duration':duration, 'ymd':ymd}

Pyspark 代码:

>>> on_time=udf(on_time, MapType(StringType(),StringType()))
>>> df_ts_list.withColumn("one_tag", on_time("ts_list")).select("one_tag").show()

错误:

Caused by: org.apache.spark.api.python.PythonException: Traceback (most recent call last):
  File "/usr/lib/spark/python/pyspark/worker.py", line 172, in main
    process()
  File "/usr/lib/spark/python/pyspark/worker.py", line 167, in process
    serializer.dump_stream(func(split_index, iterator), outfile)
  File "/usr/lib/spark/python/pyspark/worker.py", line 106, in <lambda>
    func = lambda _, it: map(mapper, it)
  File "/usr/lib/spark/python/pyspark/worker.py", line 92, in <lambda>
    mapper = lambda a: udf(*a)
  File "/usr/lib/spark/python/pyspark/worker.py", line 70, in <lambda>
    return lambda *a: f(*a)
  File "<stdin>", line 27, in on_time
  File "/usr/lib/spark/python/pyspark/sql/functions.py", line 39, in _
    jc = getattr(sc._jvm.functions, name)(col._jc if isinstance(col, Column) else col)
AttributeError: 'NoneType' object has no attribute '_jvm'

任何帮助将不胜感激!

【问题讨论】:

    标签: apache-spark ipython pyspark apache-spark-sql pyspark-sql


    【解决方案1】:

    Mariusz 的回答并没有真正帮助我。因此,如果您喜欢我发现这个,因为它是 google 上的唯一结果,并且您是 pyspark 的新手(以及一般的 spark),这对我有用。

    在我的情况下,我遇到了这个错误,因为我试图在 pyspark 环境设置之前执行 pyspark 代码。

    确保 pyspark 可用并在进行依赖于 pyspark.sql.functions 的调用之前进行设置为我解决了这个问题。

    【讨论】:

    • 作为其他人的补充......当我的 spark 会话尚未设置并且我使用装饰器定义了一个 pyspark UDF 来添加架构时,我遇到了这个错误。我通常在我的 main 中设置 spark 会话,但在这种情况下,当传递一个复杂的模式时,需要在脚本的顶部设置它。感谢您的快速提示!
    • 补充一点,在函数的默认值中使用 spark 函数时出现此错误,因为它们是在导入时评估的,而不是在调用时评估的。例如。 def func(is_test = lit(False))
    • 或者,对于像我这样愚蠢的其他人,如果您在 pandas_udf 中编写 pyspark 代码(应该接收 pandas 代码...),您可能会遇到此错误。
    【解决方案2】:

    错误消息说,在 udf 的第 27 行中,您正在调用一些 pyspark sql 函数。它与abs() 一致,所以我想你上面的某个地方调用from pyspark.sql.functions import *,它会覆盖python 的abs() 函数。

    【讨论】:

    • 有没有办法在不删除from pyspark.sql.functions import *这一行的情况下使用原来的abs()函数?
    • @mufmuf 当然,你可以使用__builtin__.abs 作为指向python函数的指针
    • 或者你可以将pyspark.sql.functions 导入为F 并使用F.function_name 调用pyspark 函数
    • 非常感谢,我没有abs,而是round
    • 我改用math.fabs()
    【解决方案3】:

    需要明确的是,很多人遇到的问题都源于一种糟糕的编程风格。那是from blah import *

    当你们这样做时

    from pyspark.sql.functions import *
    

    你覆盖了 很多 python 内置函数。我强烈推荐像

    这样的导入函数
    import pyspark.sql.functions as f
    # or 
    import pyspark.sql.functions as pyf
    

    【讨论】:

    • 这个建议帮助我改正了我在导入时使用“*”的坏习惯。希望其他人也能纠正这个问题
    【解决方案4】:

    确保您正在初始化 Spark 上下文。例如:

    spark = SparkSession \
        .builder \
        .appName("myApp") \
        .config("...") \
        .getOrCreate()
    sqlContext = SQLContext(spark)
    productData = sqlContext.read.format("com.mongodb.spark.sql").load()
    

    或者像

    spark = SparkSession.builder.appName('company').getOrCreate()
    sqlContext = SQLContext(spark)
    productData = sqlContext.read.format("csv").option("delimiter", ",") \
        .option("quote", "\"").option("escape", "\"") \
        .option("header", "true").option("inferSchema", "true") \
        .load("/path/thecsv.csv")
    

    【讨论】:

      【解决方案5】:

      udf 无法处理None 值时也会出现此异常。 例如以下代码导致相同的异常:

      get_datetime = udf(lambda ts: to_timestamp(ts), DateType())
      df = df.withColumn("datetime", get_datetime("ts"))
      

      但是这个没有:

      get_datetime = udf(lambda ts: to_timestamp(ts) if ts is not None else None, DateType())
      df = df.withColumn("datetime", get_datetime("ts"))
      

      【讨论】:

        【解决方案6】:

        我在我的 jupyter 笔记本中发现了这个错误。我添加了以下命令 导入 findspark findspark.init() sc = pyspark.SparkContext(appName="")

        它奏效了。火花上下文未准备好或停止的同样问题。

        【讨论】:

          【解决方案7】:

          我遇到了同样的问题,当我的代码中有 python 的 round() 函数时,就像@Mariusz 说的那样python 的 round() 函数被覆盖了

          解决方法是使用__builtin__.round() 而不是round(),就像@Mariusz 在他的回答中的cmets 中提到的那样。

          【讨论】:

            猜你喜欢
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 2021-11-07
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            相关资源
            最近更新 更多