【问题标题】:Pyarrow error: while running a pandas udf in pysparkPyarrow 错误:在 pyspark 中运行 pandas udf 时
【发布时间】:2020-07-04 13:43:55
【问题描述】:

我正在使用 AWS EMR (5.29) 运行 pyspark 作业,但是当我 apply a pandas udf 时收到此错误。

pyarrow.lib.ArrowInvalid: Input object was not a NumPy array

这是复制问题的虚拟代码。

import pyspark.sql.functions as F
from pyspark.sql.types import *

df = spark.createDataFrame([
    (1, "A", "X1"),
    (2, "B", "X2"),
    (3, "B", "X3"),
    (1, "B", "X3"),
    (2, "C", "X2"),
    (3, "C", "X2"),
    (1, "C", "X1"),
    (1, "B", "X1"),
], ["id", "type", "code"])

这是虚拟 udf

schema = StructType([
    StructField("code", StringType()),
])


@F.pandas_udf(schema, F.PandasUDFType.GROUPED_MAP)
def dummy_udaf(pdf):
    pdf = pdf[['code']]
    return pdf

当我运行这条线时,

df.groupby('type').apply(dummy_udaf).show()

我收到此错误:

An error occurred while calling o149.showString.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 64 in stage 12.0 failed 4 times, most recent failure: Lost task 64.3 in stage 12.0 (TID 66, ip-10-161-108-245.vpc.internal, executor 2): org.apache.spark.api.python.PythonException: Traceback (most recent call last):
  File "/mnt2/yarn/usercache/livy/appcache/application_1585015669438_0003/container_1585015669438_0003_01_000253/pyspark.zip/pyspark/worker.py", line 377, in main
    process()
  File "/mnt2/yarn/usercache/livy/appcache/application_1585015669438_0003/container_1585015669438_0003_01_000253/pyspark.zip/pyspark/worker.py", line 372, in process
    serializer.dump_stream(func(split_index, iterator), outfile)
  File "/mnt2/yarn/usercache/livy/appcache/application_1585015669438_0003/container_1585015669438_0003_01_000253/pyspark.zip/pyspark/serializers.py", line 287, in dump_stream
    batch = _create_batch(series, self._timezone)
  File "/mnt2/yarn/usercache/livy/appcache/application_1585015669438_0003/container_1585015669438_0003_01_000253/pyspark.zip/pyspark/serializers.py", line 256, in _create_batch
    arrs = [create_array(s, t) for s, t in series]
  File "/mnt2/yarn/usercache/livy/appcache/application_1585015669438_0003/container_1585015669438_0003_01_000253/pyspark.zip/pyspark/serializers.py", line 256, in <listcomp>
    arrs = [create_array(s, t) for s, t in series]
  File "/mnt2/yarn/usercache/livy/appcache/application_1585015669438_0003/container_1585015669438_0003_01_000253/pyspark.zip/pyspark/serializers.py", line 254, in create_array
    return pa.Array.from_pandas(s, mask=mask, type=t, safe=False)
  File "pyarrow/array.pxi", line 755, in pyarrow.lib.Array.from_pandas
  File "pyarrow/array.pxi", line 265, in pyarrow.lib.array
  File "pyarrow/array.pxi", line 80, in pyarrow.lib._ndarray_to_array
  File "pyarrow/error.pxi", line 84, in pyarrow.lib.check_status
pyarrow.lib.ArrowInvalid: Input object was not a NumPy array

尝试按照 here 的建议使用降级版本的箭头,并按照 here 的建议禁用 pyarrow 优化,但没有任何效果。

【问题讨论】:

  • 这在我的本地设置上运行良好。我已经安装了 pyarrow 0.13.0。另外,我有 pandas 0.22.0 和 numpy 1.17.4
  • 嘿 @Bitswazsky 是的,我在本地也可以正常工作。这仅在 AWS EMR 上发生。不知道为什么。

标签: python pandas apache-spark pyspark apache-spark-sql


【解决方案1】:

我遇到了和你一样的问题(在 AWS EMR 上)并且能够通过安装 pyarrow==0.14.1 来解决它。我不知道为什么它对您不起作用,但一种猜测是您需要在bootstrap script 中执行此安装,以便它发生在集群中的所有机器上。仅在您正在使用的笔记本中设置环境变量是不够的。希望这对您有所帮助!

【讨论】:

  • 是的,它对我不起作用,因为我在引导脚本中有另一个包,我猜,其中一个包再次升级 pyarrow。但是,当我删除所有 pip 依赖项并仅将 pyarrow==0.14.1 添加到 boostrap 脚本时,它起作用了。谢谢你的回答,我本来打算自己回答的。
猜你喜欢
  • 1970-01-01
  • 2021-07-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-03-31
  • 1970-01-01
  • 2021-08-13
相关资源
最近更新 更多