如何将PySpark DataFrame传入Java函数的Dataset参数?报错求助
问题:PySpark调用Java函数传入DataFrame时报错'_get_object_id'不存在
在PySpark Python代码中调用Java函数并传入DataFrame作为参数时,触发如下错误:
'DataFrame' object has no attribute '_get_object_id'
原Python代码
def call_java_method(spark, sc, df): df = spark.read \ .format("csv") \ .option("header", "false") \ .option("delimiter", "|").option("ignoreLeadingWhiteSpace", 'true') \ .option("ignoreTrailingWhiteSpace", 'true').load("resources/data_success.txt") df1 = df.withColumn("idx", monotonically_increasing_id()) df1.createOrReplaceTempView('databas') df1.show() df1.persist(StorageLevel.MEMORY_AND_DISK) java_import(sc._jvm, "en.test.data") func = sc._jvm.data() func.add(spark._jsparkSession, df1, 1, [2])
期望调用的Java方法
public void add(SparkSession spark, Dataset<Row> baseDataset, int count, List<Integer> keys) throws InputDataException {...}
错误原因
PySpark的Python DataFrame是Java Dataset<Row>的封装类,不能直接作为参数传给Java方法,需要获取其底层的Java对象;同时Python原生列表也无法直接映射为Java的List类型,需要转换为对应的Java集合对象。
修正后的Python代码
def call_java_method(spark, sc, df): df = spark.read \ .format("csv") \ .option("header", "false") \ .option("delimiter", "|").option("ignoreLeadingWhiteSpace", 'true') \ .option("ignoreTrailingWhiteSpace", 'true').load("resources/data_success.txt") df1 = df.withColumn("idx", monotonically_increasing_id()) df1.createOrReplaceTempView('databas') df1.show() df1.persist(StorageLevel.MEMORY_AND_DISK) java_import(sc._jvm, "en.test.data") func = sc._jvm.data() # 转换Python DataFrame为Java Dataset<Row> j_df = df1._jdf # 转换Python列表为Java ArrayList j_keys = sc._jvm.java.util.ArrayList([2]) # 调用Java方法 func.add(spark._jsparkSession, j_df, 1, j_keys)
注意事项
- 确保Java类
en.test.data所在的JAR包已添加到Spark的classpath中,否则会出现类找不到的异常 - 确认
spark._jsparkSession能正确获取到Java版的SparkSession实例 - 如果Java方法有返回值,若返回的是Java对象,需要根据实际情况转换为对应的Python对象(比如返回
Dataset<Row>的话,可以用spark.createDataFrame(j_result)转为Python DataFrame)
内容的提问来源于stack exchange,提问作者Vinay
相关产品推荐
相关产品推荐

