【问题标题】:pyspark py4j create java string arraypyspark py4j 创建java字符串数组
【发布时间】:2020-04-15 14:00:19
【问题描述】:

我想从 Pyspark 调用 FAT jar 的 main 方法。

这里是jar(Scala)的main方法的入口点:

object Main {
  def main(args: Array[String]): Unit = {
    // codes
  }
}

为了调用上述方法,我需要使用 pyspark 的 py4j 创建一个Array[String]

str_array = sc._jvm.java.lang.reflect.Array.newInstance(sc._jvm.java.lang.String, 3)
str_array[0] = "228"

loaded_class = sc._jvm.java.lang.Thread.currentThread().getContextClassLoader().loadClass("com.mycompany.Main")
loaded_class.main(str_array)

这是我得到的错误:

Py4JError: java.lang.String._get_object_id 在 JVM 中不存在

使用普通的 Py4j,我可以使用以下方法创建字符串数组:

from py4j.java_gateway import JavaGateway
gateway = JavaGateway()
gateway.new_array(gateway.jvm.java.lang.String, 4)

我尝试将对象数组传递给 main,但没有成功:

ob = sc._jvm.java.lang.Object()
ob_array = sc._jvm.java.lang.reflect.Array.newInstance(ob.getClass(), 3)
ob_array[0] = "228"
loaded_class = sc._jvm.java.lang.Thread.currentThread().getContextClassLoader().loadClass("com.mycompany.Main")
loaded_class.main(ob_array)

因错误而失败:

Py4JError: An error occurred while calling o516.main. Trace:
py4j.Py4JException: Method main([class [Ljava.lang.Object;]) does not exist
    at py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:318)
    at py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:326)
    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:238)
    at java.lang.Thread.run(Thread.java:748)

如何创建字符串数组来调用 PySpark 中的 main 方法?

【问题讨论】:

    标签: apache-spark pyspark py4j


    【解决方案1】:

    要将一些 python 字符串列表 (arr) 转换为 pyspark 中的 java 字符串数组,您可以使用下面的函数 toJStringArray

    def toJStringArray(arr):
        jarr = sc._gateway.new_array(sc._jvm.java.lang.String, len(arr))
        for i in range(len(arr)):
            jarr[i] = arr[i]
        return jarr
    
    # Usage example:
    java_arr = toJStringArray(['string1'])
    

    我在 azure databricks 中运行此代码。所以sc 是 spark 上下文的预定义全局变量。

    【讨论】:

      【解决方案2】:

      与普通的 Py4j 类似,您可以尝试使用以下方法创建它吗:

      str_array = sc._jvm._gateway.new_array(sc._jvm.java.lang.String, 4)
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-07-30
        • 2018-08-10
        • 1970-01-01
        • 2021-05-25
        • 2020-05-03
        • 1970-01-01
        相关资源
        最近更新 更多