【发布时间】:2020-07-04 16:32:09
【问题描述】:
我在 PySpark 中注册了一个 java 类,我试图将一个广播变量从 PySpark 传递给这个类中的一个方法。像这样:
from py4j.java_gateway import java_import
java_import(spark.sparkContext._jvm, "net.a.b.c.MyClass")
myPythonGateway = spark.sparkContext._jvm.MyClass()
with open("tests/fixtures/file.txt", "rb") as binary_file:
data = spark.sparkContext.broadcast(binary_file.read())
myPythonGateway.setData(data)
但这是在扔:
AttributeError:“广播”对象没有属性“_get_object_id”
但是,如果我直接传递 byte[],而不将其包装在 broadcast() 中,它可以正常工作。但是我需要广播这个变量,因为它会被重复使用。
【问题讨论】:
标签: apache-spark pyspark py4j